Beam Pipeline中CheckStopReadingFn返回true时抛出IllegalStateException问题
Apache Beam KafkaIO CheckStopReadingFn触发OffsetTracker空指针异常的原因及解决方案
问题诱因
- Apache Beam 2.29.0版本的KafkaIO存在已知设计缺陷:当
CheckStopReadingFn在分片首次调用就返回true时,此时还没有任何数据被消费,OffsetTracker的lastAttemptedOffset还未被初始化、保持null值,但当前分片的offset范围是非空的,就会触发checkDone方法的非空校验报错。 - 你自定义的
CheckStopReadingFn存在代码错误,导致逻辑误判提前返回停止信号:- 构造函数赋值错误:类成员变量为
subdirectory,但构造函数中错误将入参赋值给不存在的viewName变量,导致subdirectory始终为null,GCSUtility.filesExist的判断逻辑不符合预期,提前返回true。 - 代码语法错误:
apply方法中return = GCSUtility.filesExist(...)多写了等于号,属于语法错误(如果是粘贴手误可忽略)。
- 构造函数赋值错误:类成员变量为
修复方案
- 先修复自定义函数的代码问题:
private static class CheckStopReadingFn implements SerializableFunction<TopicPartition, Boolean> { final private String projectId; final private String bucketName; final private String subdirectory; CheckStopReadingFn(String projectId, String bucketName, String subdirectory ) { this.projectId = projectId; this.bucketName = bucketName; this.subdirectory = subdirectory; // 修正赋值错误 } @Override public Boolean apply(TopicPartition topicPartition) { return GCSUtility.filesExist(projectId, bucketName, subdirectory, Set.of(topicPartition.toString())); // 修正语法错误 } }
- 修复版本缺陷:
- 优先升级Apache Beam版本到2.34.0及以上,该版本官方已修复
CheckStopReadingFn提前触发停止时的空指针异常问题。 - 若暂时无法升级版本,可在
CheckStopReadingFn逻辑中增加兼容处理:确保对应分区至少消费过一条数据、或者分区的起始偏移量等于最新偏移量时,再返回true,避免未消费任何数据就触发停止的场景。
- 优先升级Apache Beam版本到2.34.0及以上,该版本官方已修复
内容的提问来源于stack exchange,提问作者Fabio
相关产品推荐
相关产品推荐

