Apache Flink 1.4.2对接Kafka源时JSON字符串被截断的问题求助
我之前在维护基于Flink 1.4.x的项目时也碰到过一模一样的问题——大JSON被截断到4095字符导致解析失败,结合你给出的版本信息,这里给你拆解原因和可行的解决办法:
为什么会被截断到4095字符?
Flink 1.4.2集成的是Kafka 0.10.x版本的客户端,这个版本的Kafka Consumer有个默认配置max.partition.fetch.bytes,值为4096字节。这个配置控制的是Consumer每次从单个Kafka分区拉取的总字节数,其中还要包含Kafka的消息协议头(大概1字节左右),所以实际能拿到的消息内容就被限制在了4095字节。当你的JSON字符串超过这个长度时,就会被截断,剩下的内容会被当成下一条“不完整的消息”,自然会导致JSON解析失败。
而且这个问题很容易被忽略,因为Kafka Broker端的默认max.message.bytes是1MB,本身是允许大消息的,但Consumer端的拉取限制没跟上,就会出现截断。
解决方案
针对你的情况,不需要修改JSON的大小,只需要调整两处配置就能解决:
1. 调整Flink Kafka Consumer的拉取字节限制
在创建Flink Kafka Consumer的时候,通过Properties手动设置max.partition.fetch.bytes为你实际需要的大小(比如你的JSON最大是2MB,就设为210241024=2097152)。示例代码如下:
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer010; import org.apache.flink.api.common.serialization.SimpleStringSchema; import java.util.Properties; // 初始化Kafka配置 Properties kafkaProps = new Properties(); kafkaProps.put("bootstrap.servers", "你的Kafka Broker地址"); kafkaProps.put("group.id", "你的消费组ID"); // 设置拉取的最大字节数,根据你的JSON实际大小调整 kafkaProps.put("max.partition.fetch.bytes", "2097152"); // 创建Consumer FlinkKafkaConsumer010<String> kafkaConsumer = new FlinkKafkaConsumer010<>( "你的Kafka Topic名称", new SimpleStringSchema(), // 你用的字符串反序列化 kafkaProps );
2. 同步调整Kafka Broker的消息大小限制(如果需要)
如果你的JSON大小超过了Kafka Broker默认的max.message.bytes(默认1MB),还需要修改Kafka Broker的server.properties配置:
# 设置允许的最大单条消息大小,要大于等于上面设置的max.partition.fetch.bytes max.message.bytes=2097152
修改后记得重启Kafka Broker生效。
额外提醒
因为你用的Flink 1.4.2是比较老的版本,后续的Flink版本(比如1.10+)已经调整了这些默认配置,并且对Kafka大消息的处理做了优化。如果后续有机会升级Flink版本,也能避免这类问题,但当前版本下,调整上述两个配置是最直接有效的办法。
内容的提问来源于stack exchange,提问作者ankita.gulati

