多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
相关产品推荐
相关产品推荐

