Python Apache Beam 2.32.0 ReadFromKafka消费Kafka数据报错
问题原因及解决方案
第一个错误(无法编码null字节数组)原因
Kafka控制台生产者发送消息时默认不设置消息key,key值为null。Apache Beam默认使用ByteArrayCoder处理ReadFromKafka返回的KafkaRecord的key和value,而ByteArrayCoder不支持编码null值,因此抛出编码异常。
第二个错误(无法构建Kafka IO转换)原因
你混淆了序列化器和反序列化器的使用场景:
Serializer是Kafka生产者发送消息时用的序列化类ReadFromKafka是消费逻辑,需要传入Deserializer结尾的反序列化类
你传入的StringSerializer属于生产者序列化类,类型不匹配导致Beam无法构造对应的跨语言IO转换。此外如果反序列化器输出类型和默认的ByteArrayCoder不匹配,也会触发该类错误。
正确配置方案
提供两种可选的修复方式:
方案1:显式配置反序列化器和对应Coder
指定字符串类型的反序列化器,同时配套使用字符串Coder:
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.io.external.kafka import ReadFromKafka pipeline_options = PipelineOptions(["--runner=DirectRunner"]) def run(): with beam.Pipeline(options=pipeline_options) as p: _ = ( p | 'ReadData' >> ReadFromKafka( consumer_config={"bootstrap.servers": "localhost:9092"}, topics=["my-first-topic"], key_deserializer="org.apache.kafka.common.serialization.StringDeserializer", value_deserializer="org.apache.kafka.common.serialization.StringDeserializer", key_coder=beam.coders.StrUtf8Coder(), value_coder=beam.coders.StrUtf8Coder() ) | 'PrintData' >> beam.Map(print) ) if __name__ == "__main__": run()
方案2:默认配置下自行处理null值
不修改反序列化器配置,在消费逻辑中提前处理null的key/value,避免Coder编码null:
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.io.external.kafka import ReadFromKafka pipeline_options = PipelineOptions(["--runner=DirectRunner"]) def run(): with beam.Pipeline(options=pipeline_options) as p: _ = ( p | 'ReadData' >> ReadFromKafka( consumer_config={"bootstrap.servers": "localhost:9092"}, topics=["my-first-topic"], ) | 'HandleNull' >> beam.Map(lambda record: ( record.key.decode('utf-8') if record.key is not None else '', record.value.decode('utf-8') if record.value is not None else '' )) | 'PrintData' >> beam.Map(print) ) if __name__ == "__main__": run()
内容的提问来源于stack exchange,提问作者user3595632
相关产品推荐
相关产品推荐

