You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.04 11:06:02