Flink Python Datastream API Kafka Producer Sink序列化报错如何解决
问题根因
- 类型不匹配报错:报错堆栈最后一行的
java.lang.ClassCastException: class [B cannot be cast to class java.lang.String是核心错误,[B是Java中字节数组的标识。PyFlink的map算子如果没有显式指定输出类型,会默认将Python侧返回的字符串序列化为字节数组传递到Java侧,而SimpleStringSchema的serialize方法要求接收String类型参数,收到字节数组就会抛出类型转换异常。 - 多余的反序列化逻辑:你自定义的
my_map函数中连续调用了两次json.loads,除非你的源数据是经过两次JSON序列化的嵌套字符串,否则该逻辑会直接抛出JSON解析异常。
修复方案
import json from pyflink.common.typeinfo import Types from pyflink.common.serialization import SimpleStringSchema from pyflink.datastream import StreamExecutionEnvironment from pyflink.datastream.connectors import FlinkKafkaConsumer, FlinkKafkaProducer def my_map(obj): # 仅保留一次反序列化逻辑,适配标准JSON格式的源数据 json_obj = json.loads(obj) return json.dumps(json_obj["name"]) def datastream_api_demo(): env = StreamExecutionEnvironment.get_execution_environment() env.add_jars("file:///Users/niaz/Downloads/f2.jar") deserialization_schema = SimpleStringSchema() kafka_source = FlinkKafkaConsumer( topics='test_source_topic_input', deserialization_schema=deserialization_schema, properties={'bootstrap.servers': 'localhost:9092', 'group.id': 'test_group'}) ds = env.add_source(kafka_source) # 新增output_type参数,显式指定算子输出为字符串类型 ds = ds.map(lambda a: my_map(a), output_type=Types.STRING()) serialization_schema = SimpleStringSchema() kafka_sink = FlinkKafkaProducer( topic='test_sink_topic_4', serialization_schema=serialization_schema, producer_config={'bootstrap.servers': 'localhost:9092', 'group.id': 'test_group'}) ds.add_sink(kafka_sink) env.execute('datastream_api_demo') if __name__ == '__main__': datastream_api_demo()
内容的提问来源于stack exchange,提问作者Mujahid Niaz
相关产品推荐
相关产品推荐

