You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 15:50:22