EventHub未返回作业完成状态,任务进入轮询分支求助
问题排查与疑问解答
问题描述
我尝试获取Job状态,但Task始终无法完成,一直进入else轮询块,无法定位问题所在。
实现代码
主逻辑代码
#region Job Status Action EventProcessorClient processorClient = null; BlobContainerClient storageClient = null; MediaServicesEventProcessor mediaEventProcessor = null; //log.LogInformation($"Job '{jobName}' submitted."); try { Console.WriteLine("Creating a new client to process events from an Event Hub..."); var credential = new DefaultAzureCredential(); var EventHubConnectionString = "Endpoint=sb://komurojusaikrishn.servicebus.windows.net/;SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey=<Shared Access key>"; //var storageConnectionString = string.Format("DefaultEndpointsProtocol=https;AccountName={0};AccountKey={1}", // config.AccountName, config.StorageAccountKey); var storageConnectionString = string.Format("DefaultEndpointsProtocol=https;AccountName=<Account Name>;AccountKey=<Account key>;EndpointSuffix=core.windows.net"); var blobContainerName = "asset-fcc961b8-cc25-4f4f-8617-9c767d865572";//"blob-01182023";//config.StorageContainerName; var eventHubsConnectionString = EventHubConnectionString; var eventHubName = "<Event Hub Name>"; var consumerGroup = "$Default";//config.EventHubConsumerGroup; storageClient = new BlobContainerClient( storageConnectionString, blobContainerName); processorClient = new EventProcessorClient( storageClient, consumerGroup, eventHubsConnectionString, eventHubName); AutoResetEvent jobWaitingEvent = new AutoResetEvent(false); IList<Task> tasks = new List<Task>(); Task jobTask = Task.Run(() => jobWaitingEvent.WaitOne()); tasks.Add(jobTask); // 30 minutes timeout. var cancellationSource = new CancellationTokenSource(); var timeout = Task.Delay(5 * 60 * 1000, cancellationSource.Token); tasks.Add(timeout); mediaEventProcessor = new MediaServicesEventProcessor(jobName, jobWaitingEvent, null); processorClient.ProcessEventAsync += mediaEventProcessor.ProcessEventsAsync; processorClient.ProcessErrorAsync += mediaEventProcessor.ProcessErrorAsync; await processorClient.StartProcessingAsync(cancellationSource.Token); //**此处条件始终不满足,因为jobTask未完成,进入else块** if (await Task.WhenAny(tasks) == jobTask) { cancellationSource.Cancel(); job = await client.Jobs.GetAsync(config.ResourceGroup, config.AccountName, ObjReqconfig.TransformName, jobName); } else { jobWaitingEvent.Set(); throw new Exception("Timeout occurred."); } } catch (Exception e) { Console.WriteLine("Warning: Failed to connect to Event Hub, please refer README for Event Hub and storage settings."); Console.WriteLine(e.Message); Console.WriteLine("Polling job status..."); job = await JobStatusChecker.WaitForJobToFinishAsync(client, config.ResourceGroup, config.AccountName, ObjReqconfig.TransformName, jobName); } finally { if (processorClient != null) { Console.WriteLine("Job final state received, Stopping the event processor..."); await processorClient.StopProcessingAsync(); Console.WriteLine(); processorClient.ProcessEventAsync -= mediaEventProcessor.ProcessEventsAsync; processorClient.ProcessErrorAsync -= mediaEventProcessor.ProcessErrorAsync; } } #endregion //TimeSpan elapsed = DateTime.Now - startedTime;
MediaServicesEventProcessor.cs代码
public class MediaServicesEventProcessor { private readonly AutoResetEvent jobWaitingEvent; private readonly string jobName; private readonly string liveEventName; public MediaServicesEventProcessor(string jobName, AutoResetEvent jobWaitingEvent, string liveEventName) { this.jobName = jobName; this.jobWaitingEvent = jobWaitingEvent; this.liveEventName = liveEventName; } public Task ProcessErrorAsync(ProcessErrorEventArgs args) { try { Console.WriteLine("Error in the EventProcessorClient"); Console.WriteLine($"\tOperation: { args.Operation }"); Console.WriteLine($"\tException: { args.Exception }"); Console.WriteLine(""); } catch { } return Task.CompletedTask; } public Task ProcessEventsAsync(ProcessEventArgs args) { if (args.HasEvent) { PrintJobEvent(args.Data); } return args.UpdateCheckpointAsync(); } private void PrintJobEvent(Azure.Messaging.EventHubs.EventData eventData) { // data = Encoding.UTF8.GetString(eventData.EventBody); var options = new JsonSerializerOptions { PropertyNameCaseInsensitive = true }; //EventGridEvent[] egEvents var jsonArray = JsonSerializer.Deserialize<EventGridEvent[]>(eventData.EventBody.ToString(), options); //foreach (Azure.Messaging.EventGrid.EventGridEvent e in jsonArray) foreach (EventGridEvent e in jsonArray) { var subject = e.Subject; var topic = e.Topic; var eventType = e.EventType; var eventTime = e.EventTime; string eventSourceName = Regex.Replace(subject, @"^.*/", ""); if (eventSourceName != jobName && eventSourceName != liveEventName) { return; } // Log the time and type of event Console.WriteLine($"{eventTime} EventType: {eventType}"); switch (eventType) { // Job state change events case "Microsoft.Media.JobStateChange": case "Microsoft.Media.JobScheduled": case "Microsoft.Media.JobProcessing": case "Microsoft.Media.JobCanceling": case "Microsoft.Media.JobFinished": case "Microsoft.Media.JobCanceled": case "Microsoft.Media.JobErrored": { MediaJobStateChangeEventData jobEventData = JsonSerializer.Deserialize<MediaJobStateChangeEventData>(e.Data.ToString(), options); Console.WriteLine($"Job state changed for JobId: {eventSourceName} PreviousState: {jobEventData.PreviousState} State: {jobEventData.State}"); // For final states, send a message to notify that the job has finished. if (eventType == "Microsoft.Media.JobFinished" || eventType == "Microsoft.Media.JobCanceled" || eventType == "Microsoft.Media.JobErrored") { // Job finished, send a message. if (jobWaitingEvent != null) { jobWaitingEvent.Set(); } } } break; // Job output state change events case "Microsoft.Media.JobOutputStateChange": case "Microsoft.Media.JobOutputScheduled": case "Microsoft.Media.JobOutputProcessing": case "Microsoft.Media.JobOutputCanceling": case "Microsoft.Media.JobOutputFinished": case "Microsoft.Media.JobOutputCanceled": case "Microsoft.Media.JobOutputErrored": { MediaJobOutputStateChangeEventData jobEventData = JsonSerializer.Deserialize<MediaJobOutputStateChangeEventData>(e.Data.ToString(), options); Console.WriteLine($"Job output state changed for JobId: {eventSourceName} PreviousState: {jobEventData.PreviousState} " + $"State: {jobEventData.Output.State} Progress: {jobEventData.Output.Progress}%"); } break; // Job output progress event case "Microsoft.Media.JobOutputProgress": { MediaJobOutputProgressEventData jobEventData = JsonSerializer.Deserialize<MediaJobOutputProgressEventData>(e.Data.ToString(), options); Console.WriteLine($"Job output progress changed for JobId: {eventSourceName} Progress: {jobEventData.Progress}%"); } break; // LiveEvent Stream-level events // See the following documentation for updated schemas - https://docs.microsoft.com/azure/media-services/latest/monitoring/media-services-event-schemas#live-event-types case "Microsoft.Media.LiveEventConnectionRejected": { MediaLiveEventConnectionRejectedEventData liveEventData = JsonSerializer.Deserialize<MediaLiveEventConnectionRejectedEventData>(e.Data.ToString(), options); Console.WriteLine($"LiveEvent connection rejected. IngestUrl: {liveEventData.IngestUrl} StreamId: {liveEventData.StreamId} " + $"EncoderIp: {liveEventData.EncoderIp} EncoderPort: {liveEventData.EncoderPort}"); } break; case "Microsoft.Media.LiveEventEncoderConnected": { MediaLiveEventEncoderConnectedEventData liveEventData = JsonSerializer.Deserialize<MediaLiveEventEncoderConnectedEventData>(e.Data.ToString(), options); Console.WriteLine($"LiveEvent encoder connected. IngestUrl: {liveEventData.IngestUrl} StreamId: {liveEventData.StreamId} " + $"EncoderIp: {liveEventData.EncoderIp} EncoderPort: {liveEventData.EncoderPort}"); } break; case "Microsoft.Media.LiveEventEncoderDisconnected": { MediaLiveEventEncoderDisconnectedEventData liveEventData = JsonSerializer.Deserialize<MediaLiveEventEncoderDisconnectedEventData>(e.Data.ToString(), options); Console.WriteLine($"LiveEvent encoder disconnected. IngestUrl: {liveEventData.IngestUrl} StreamId: {liveEventData.StreamId} " + $"EncoderIp: {liveEventData.EncoderIp} EncoderPort: {liveEventData.EncoderPort}"); } break; case "Microsoft.Media.LiveEventIncomingDataChunkDropped": { MediaLiveEventIncomingDataChunkDroppedEventData liveEventData = JsonSerializer.Deserialize<MediaLiveEventIncomingDataChunkDroppedEventData>(e.Data.ToString(), options); Console.WriteLine($"LiveEvent data chunk dropped. LiveEventId: {eventSourceName} ResultCode: {liveEventData.ResultCode}"); } break; case "Microsoft.Media.LiveEventIncomingStreamReceived": { MediaLiveEventIncomingStreamReceivedEventData liveEventData = JsonSerializer.Deserialize<MediaLiveEventIncomingStreamReceivedEventData>(e.Data.ToString(), options); Console.WriteLine($"LiveEvent incoming stream received. IngestUrl: {liveEventData.IngestUrl} EncoderIp: {liveEventData.EncoderIp} " + $"EncoderPort: {liveEventData.EncoderPort}"); } break; case "Microsoft.Media.LiveEventIncomingStreamsOutOfSync": { //MediaLiveEventIncomingStreamsOutOfSyncEventData eventData = JsonSerializer.Deserialize<MediaLiveEventIncomingStreamsOutOfSyncEventData>(e.Data.ToString(), options);; Console.WriteLine($"LiveEvent incoming audio and video streams are out of sync. LiveEventId: {eventSourceName}"); } break; case "Microsoft.Media.LiveEventIncomingVideoStreamsOutOfSync": { //MediaLiveEventIncomingVideoStreamsOutOfSyncEventData eventData =JsonSerializer.Deserialize<MediaLiveEventIncomingVideoStreamsOutOfSyncEventData>(e.Data.ToString(), options);; Console.WriteLine($"LeveEvent incoming video streams are out of sync. LiveEventId: {eventSourceName}"); } break; case "Microsoft.Media.LiveEventIngestHeartbeat": { MediaLiveEventIngestHeartbeatEventData liveEventData = JsonSerializer.Deserialize<MediaLiveEventIngestHeartbeatEventData>(e.Data.ToString(), options); Console.WriteLine($"LiveEvent ingest heart beat. TrackType: {liveEventData.TrackType} State: {liveEventData.State} Healthy: {liveEventData.Healthy}"); } break; case "Microsoft.Media.LiveEventTrackDiscontinuityDetected": { MediaLiveEventTrackDiscontinuityDetectedEventData liveEventData = JsonSerializer.Deserialize<MediaLiveEventTrackDiscontinuityDetectedEventData>(e.Data.ToString(), options); Console.WriteLine($"LiveEvent discontinuity in the incoming track detected. LiveEventId: {eventSourceName} TrackType: {liveEventData.TrackType} " + $"Discontinuity gap: {liveEventData.DiscontinuityGap}"); } break; case "Microsoft.Media.LiveEventChannelArchiveHeartbeatEvent": { Console.WriteLine($"LiveEvent archive heartbeat event detected. LiveEventId: {eventSourceName}"); Console.WriteLine(e.Data.ToString()); } break; } } } }
排查要点
Event Hub事件接收验证
- 检查Event Grid是否正确配置,将Media Services的Job事件转发到指定Event Hub,确认订阅状态激活且筛选条件包含目标Job的状态变更事件。
- 验证Event Hub连接字符串、名称、消费者组的正确性,确保程序有读取权限。
- 查看
ProcessErrorAsync的输出,确认是否有连接或权限异常。
事件匹配逻辑问题
- 确认
PrintJobEvent中提取eventSourceName的正则@"^.*/"是否能正确从subject中截取Job名称,若subject格式不符会导致事件被直接忽略。 - 添加日志验证
eventSourceName和jobName是否匹配,避免因格式差异导致逻辑分支不触发。
- 确认
AutoResetEvent使用问题
- 在
jobWaitingEvent.Set()前后添加日志,确认是否执行到该代码分支。 - 检查是否有其他操作意外重置了
AutoResetEvent,导致jobTask持续等待。
- 在
任务等待逻辑问题
- 验证
processorClient.StartProcessingAsync是否成功启动,添加日志确认启动状态。 - 检查
cancellationSource是否被意外取消,导致事件处理器停止接收事件。
- 验证
Blob容器疑问解答
EventProcessorClient使用的Blob容器用于存储检查点数据,既可以复用输入Asset的容器,也可以使用独立容器:
- 复用现有容器:需确保容器有写入权限,且检查点文件(默认存储在
$default目录)不会和Asset文件冲突。 - 使用独立容器:更清晰,避免业务数据与检查点数据混合,便于管理,生产环境建议优先选择。
内容的提问来源于stack exchange,提问作者pranayrebel
相关产品推荐
相关产品推荐

