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

多Scheduler场景下,如何在出错时停止KCL的processRecord方法?

解决方案:优雅停止单个KCL Worker实例

问题根源

  • Thread.currentThread().stop()是Java已废弃的危险方法,会强制终止线程,导致AWS客户端(如Kinesis/DynamoDB客户端)资源未正常释放,进而引发Client is closed错误。
  • 直接抛出异常无法停止Worker,因为KCL默认会捕获异常并重启处理循环,以保证消费可用性。

正确实现方式

要实现单个Worker故障时仅停止自身,需持有Worker实例引用,通过KCL原生的Worker.stop()方法优雅终止,而非直接操作线程。

步骤1:重构代码,持有Worker实例

为每个Kinesis流维护独立的Worker实例,将实例存储在线程安全的映射中(比如Map<String, Worker>),键用流名称,方便定位需停止的Worker。同时在RecordProcessor中保留当前Worker的引用。

步骤2:修改processRecords方法,触发优雅停止

捕获到需终止Worker的异常时,先完成checkpoint避免重复处理,再标记停止状态,最后异步调用Worker.stop()(避免阻塞当前处理线程)。

示例代码:

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.atomic.AtomicBoolean;

// RecordProcessor类内部新增
private final AtomicBoolean shouldStop = new AtomicBoolean(false);
private Worker currentWorker;

public void setCurrentWorker(Worker worker) {
    this.currentWorker = worker;
}

@SneakyThrows
@Override
public void processRecords(ProcessRecordsInput processRecordsInput) {
    if (shouldStop.get()) {
        return; // 已标记停止,直接返回
    }

    for (KinesisClientRecord s : processRecordsInput.records()) {
        try {
            // 业务操作:Glue校验、DynamoDB存储等逻辑
        } catch (GlueSchemaException e) {
            // 处理Schema异常,不终止Worker
        } catch (Exception ex) {
            LOGGER.info("准备停止当前Worker: {}", Thread.currentThread().getName());
            // 先完成checkpoint,防止重复消费
            processRecordsInput.checkpointer().checkpoint();
            // 标记需要停止
            shouldStop.set(true);
            // 异步执行停止操作,避免阻塞当前处理流程
            CompletableFuture.runAsync(() -> {
                if (currentWorker != null) {
                    currentWorker.stop();
                    LOGGER.info("Worker已成功停止");
                }
            });
            // 终止当前批次记录处理
            break;
        }
        // 单条记录处理完成后checkpoint
        processRecordsInput.checkpointer().checkpoint();
    }
}

步骤3:启动Worker时绑定实例

创建每个流的Worker时,将实例注入对应的RecordProcessor:

import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

// 示例:批量创建并启动Worker
Map<String, Worker> workerMap = new ConcurrentHashMap<>();

for (String streamName : streamNames) {
    YourRecordProcessor processor = new YourRecordProcessor();
    Worker worker = Worker.builder()
            .recordProcessorFactory(() -> processor)
            .streamName(streamName)
            // 配置Kinesis客户端、DynamoDB客户端等参数
            .build();
    processor.setCurrentWorker(worker);
    workerMap.put(streamName, worker);
    // 启动Worker线程
    new Thread(worker::run).start();
}

关键注意事项

  • 禁止使用Thread.stop(),它会引发不可预测的资源泄漏和数据不一致问题。
  • 必须使用Worker.stop(),该方法会优雅终止KCL处理循环、释放资源,避免客户端报错。
  • 异步停止Worker是为了保证当前checkpoint操作完成,避免阻塞处理线程。
  • 每个Worker对应独立的RecordProcessor和AWS客户端实例,确保单个Worker停止不会影响其他Worker的正常运行。

内容的提问来源于stack exchange,提问作者Prathviraj Singh Chouhan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 07:03:16