Kafka Connect S3 Source抛出java.io.IOException流已关闭异常咨询
问题答复
1. 异常触发原因与修复方法
触发原因
该异常是confluentinc-kafka-connect-s3-source-2.1.1版本对gzip压缩对象的读取逻辑bug导致:
- 当S3 Sink端开启gzip压缩写入JSON格式对象时,旧版本S3 Source的JsonFormat解析器调用
GZIPInputStream读取内容,读到压缩数据的结束标记后就停止读取,没有继续消费底层S3ObjectInputStream中gzip格式的尾部校验字节,导致输入流存在未读的剩余字节。 - AWS SDK的S3AbortableInputStream检测到流未被完全消费就被关闭,主动输出如下警告并中断HTTP连接释放资源,后续解析流程再尝试读取该流时,就会抛出
java.io.IOException: Attempted read on closed stream异常:
WARN Not all bytes were read from the S3ObjectInputStream, aborting HTTP connection. This is likely an error and may result in sub-optimal behavior. Request only the bytes you need via a ranged GET or drain the input stream after use. (com.amazonaws.services.s3.internal.S3AbortableInputStream:178)
- 关闭Sink端gzip压缩后问题不再复现,完全匹配该bug的触发条件。
修复方案
- 方案一(已验证有效):直接关闭S3 Sink端的gzip压缩配置,即可规避该兼容问题
- 方案二(长期修复):将S3 Source连接器版本从2.1.1升级到2.1.2及以上,官方已经在新版本修复了压缩对象读取时流未正常耗尽的bug
- 方案三(临时兜底):如果既不能关闭压缩也暂时无法升级版本,可以在连接器配置中增大S3客户端的socket超时、HTTP连接超时参数,降低连接提前中断的概率,但该方案稳定性不如前两种。
2. S3 Source连接器开源属性说明
Confluent官方的Kafka Connect S3 Source连接器是开源项目,采用Confluent社区许可证开源,代码归属于Confluent维护的云存储连接器集合项目,S3 Source是该项目下的子模块,没有独立拆分的单独代码仓库。
内容的提问来源于stack exchange,提问作者Yordan Pavlov
相关产品推荐
相关产品推荐

