我有一個 Timer Triggered 函式,它將物件發送到這樣的服務總線主題
[FunctionName("ProcessTermedEmployee")]
public static async Task Run(
[TimerTrigger("%TimerInterval%"
#if DEBUG
, RunOnStartup=true
#endif
)] TimerInfo myTimer,
[ServiceBus("%TermedEmployeeTopicName%", Connection = "ServiceBusConnectionString")] IAsyncCollector<string> termedEmployeeCollector,
ILogger log)
{
log.LogInformation("ProcessTermedEmployee Timer function invoked.");
try
{
// var termedEmployees = new List<TermedEmployee>();
using (SqlConnection conn = new SqlConnection(Environment.GetEnvironmentVariable("HRSQLConNString")))
{
conn.Open();
SqlCommand cmd = new SqlCommand("hp_Cdi_Get_Termed_Employees_Recent", conn);
cmd.CommandType = CommandType.StoredProcedure;
cmd.Parameters.Add(new SqlParameter("@TermedAfterDate", DateTime.Now.AddDays(-10)));
using (SqlDataReader rdr = cmd.ExecuteReader())
{
try
{
while (rdr.Read())
{
TermedEmployee termedEmployee = ConvertTermedEmployee(rdr);
// termedEmployees.Add(termedEmployee);
await
termedEmployeeCollector.AddAsync(JsonConvert.SerializeObject(termedEmployee));
}
}
catch (Exception ex)
{
ex.Data.Add("row", rdr);
log.LogCritical(ex, ex.Message);
throw ex;
}
}
}
log.LogInformation("ProcessTermedEmployee Timer function finished.");
}
catch (Exception ex)
{
log.LogCritical(ex, ex.Message);
throw ex;
}
}
然后我有一個服務總線觸發功能,假設接收訊息并將其添加到 CosmosDB
public static class LogTermedEmployee
{
[FunctionName("LogTermedEmployee")]
public static async Task Run([ServiceBusTrigger("%TermedEmployeeTopicName%", "%ServiceBusSubscriptionName%", Connection = "ServiceBusConnectionString")] BrokeredMessage message,
[CosmosDB(
databaseName: "%CosmosDbName%",
collectionName: "%TermedEmployeesLogCollection%",
ConnectionStringSetting = "CosmosConnection")]
IAsyncCollector<TermedEmployeeLog> termedEmployeesLogCollection,
ILogger log)
{
log.LogInformation($"C# ServiceBus topic trigger function LogTermedEmployee,
processed at {DateTime.Now}");
try
{
if(message.GetBody<TermedEmployee>() != null)
{
//StreamReader reader = new StreamReader(message.Body);
// string s = System.Text.Encoding.Default.GetString(message.Body);
var termedEmp = JsonConvert.DeserializeObject<TermedEmployee>
(message.GetBody<string>());
await termedEmployeesLogCollection.AddAsync(new TermedEmployeeLog() {
TermedEmployee = termedEmp, DateOfProcess = DateTime.Now, Id =
Guid.NewGuid(), PartitionKey = "TermedEmployee", ModifiedOn =
DateTime.Now, ModifiedBy = "TermedEployeeLogService", MessageId = ""
});// message.MessageId });
}
}
catch (Exception ex)
{
log.LogCritical(ex, ex.Message);
throw ex;
}
}
<PackageReference Include="Azure.Messaging.ServiceBus" Version="7.6.0" />
<PackageReference Include="AzureFunctions.Extensions.DependencyInjection" Version="1.1.3" />
<PackageReference Include="CDI.Utilities.LogHelpers" Version="0.1.3" />
<PackageReference Include="Microsoft.AspNetCore.Mvc.Abstractions" Version="2.2.0" />
<PackageReference Include="Microsoft.Azure.Functions.Extensions" Version="1.1.0" />
<PackageReference Include="Microsoft.Azure.ServiceBus" Version="5.2.0" />
<PackageReference Include="Microsoft.Azure.WebJobs.Extensions" Version="4.0.1" />
<PackageReference Include="Microsoft.Azure.WebJobs.Extensions.CosmosDB" Version="3.0.10" />
<PackageReference Include="Microsoft.Azure.WebJobs.Extensions.Http" Version="3.0.12" />
<PackageReference Include="Microsoft.Azure.WebJobs.Extensions.ServiceBus" Version="5.2.0" />
<PackageReference Include="Microsoft.Azure.WebJobs.Script.ExtensionsMetadataGenerator" Version="4.0.1" />
<PackageReference Include="Microsoft.NET.Sdk.Functions" Version="4.0.1" />
<PackageReference Include="System.Data.SqlClient" Version="4.8.3" />
<PackageReference Include="WindowsAzure.ServiceBus" Version="6.2.2" />
問題
在 ServiceBus 觸發器的簽名中
如果我使用 BrokeredMessage 訊息 message.GetBody() 為空
如果我使用 Obsoleted Message 物件,那么 Message.Body 為空
但如果我使用簡單的字串,我會得到正確的 Json。
請問有什么建議嗎?
uj5u.com熱心網友回復:
我建議在您的觸發器中更改BrokeredMessage為ServiceBusReceivedMessagein:
public static async Task Run(
[ServiceBusTrigger(
"%TermedEmployeeTopicName%",
"%ServiceBusSubscriptionName%",
Connection = "ServiceBusConnectionString")]
ServiceBusReceivedMessage message,
[CosmosDB(
databaseName: "%CosmosDbName%",
collectionName: "%TermedEmployeesLogCollection%",
ConnectionStringSetting = "CosmosConnection")]
IAsyncCollector<TermedEmployeeLog> termedEmployeesLogCollection,
ILogger log)
我相信這將解決系結問題。
附加背景關系
從 v5.0.0 開始,該Microsoft.Azure.WebJobs.Extensions.ServiceBus包開始在Azure.Messaging.ServiceBus內部使用。新包中沒有BrokeredMessage型別。傳入訊息鍵入為ServiceBusReceivedMessage,傳出訊息鍵入為ServiceBusMessage。
可以在Microsoft.Azure.WebJobs.Extensions.ServiceBus 檔案中找到更多資訊。
轉載請註明出處,本文鏈接:https://www.uj5u.com/qianduan/426218.html
上一篇:非遞減陣列邏輯失敗
