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

Unity中调用Kafka Consumer.Close()导致程序强制退出问题求助

Unity中Confluent.Kafka消费者安全关闭解决方案

问题根源

直接在OnApplicationQuit调用_consumer.Close()时,后台消费线程可能正阻塞在_consumer.Consume(cancellationToken)方法上,主线程与消费线程的资源竞争导致Unity无响应。必须先让消费线程正常退出循环,再执行消费者的关闭操作。

修复步骤及代码

  • 保存消费任务的引用,以便后续等待其完成
  • 在OnApplicationQuit中先取消令牌,通知消费线程退出
  • 等待消费任务完全结束后,再安全关闭消费者
  • 清理CancellationTokenSource资源

修改后的完整代码:

using Confluent.Kafka;
using System;
using System.Threading;
using System.Threading.Tasks;
using UnityEngine;

public class _KafkaConsume : MonoBehaviour
{
    private IConsumer<Ignore, string> _consumer;
    private const string _topic = "testTopic";
    private CancellationTokenSource _cancellationTokenSource;
    private Task _consumeTask; // 保存消费任务引用

    private async void Start()
    {
        _cancellationTokenSource = new CancellationTokenSource();
        var config = new ConsumerConfig
        {
            BootstrapServers = "ip",
            AutoOffsetReset = AutoOffsetReset.Latest,
            GroupId = "consumerGroup",
            BrokerAddressFamily = BrokerAddressFamily.V4,
            EnableAutoCommit = false,
        };

        _consumer = new ConsumerBuilder<Ignore, string>(config).Build();
        _consumer.Subscribe(_topic);

        // 保存消费任务,避免Start方法阻塞
        _consumeTask = Task.Run(() => ConsumeMessagesAsync(_cancellationTokenSource.Token));
        await _consumeTask;
    }

    private async Task ConsumeMessagesAsync(CancellationToken cancellationToken)
    {
        while (!cancellationToken.IsCancellationRequested)
        {
            try
            {
                // 带超时的Consume调用,避免令牌取消时长时间阻塞
                var consumeResult = _consumer.Consume(TimeSpan.FromSeconds(1), cancellationToken);
                if (consumeResult != null && consumeResult.Message != null)
                {
                    Debug.Log($"Received message: {consumeResult.Message.Value}");
                }
            }
            catch (ConsumeException e)
            {
                Debug.LogError($"Error consuming message: {e.Error}");
            }
            catch (OperationCanceledException)
            {
                // 捕获令牌取消异常,正常退出循环
                break;
            }
        }
    }

    private void OnApplicationQuit()
    {
        // 1. 取消令牌,通知消费线程退出
        _cancellationTokenSource?.Cancel();

        // 2. 等待消费任务完成,最多等待3秒避免无响应
        if (_consumeTask != null && !_consumeTask.IsCompleted)
        {
            _consumeTask.Wait(TimeSpan.FromSeconds(3));
        }

        // 3. 安全关闭消费者并释放资源
        _consumer?.Close();
        _consumer?.Dispose();
        _cancellationTokenSource?.Dispose();
    }
}

关键修改说明

  • 添加_consumeTask引用:跟踪后台消费任务,确保能等待其完成
  • 带超时的Consume调用:每秒检查一次退出条件,避免长时间阻塞
  • 捕获OperationCanceledException:明确处理令牌取消情况,确保线程正常退出
  • 分步骤关闭流程:先取消令牌→等待任务结束→关闭消费者,彻底避免线程冲突

内容的提问来源于stack exchange,提问作者Seol

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 15:02:28