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

PyFlink 1.16 KafkaSink报AttributeError: 'NoneType'无startswith属性

问题场景

使用PyFlink 1.16的KafkaSource/KafkaSink实现跨Topic数据转发,读取数据功能正常且能打印结果,但KafkaSink写入时触发如下异常:

NOTE: Picked up JDK_JAVA_OPTIONS: --add-opens java.base/java.lang=ALL-UNNAMED --add-opens java.base/java.util=ALL-UNNAMED --add-opens java.base/java.util.concurrent.atomic=ALL-UNNAMED
Traceback (most recent call last):
      File "/home/.../PycharmProjects/reddit-anomaly-detection-job/main.py", line 75, in <module>
        main()
      File "/home/.../PycharmProjects/reddit-anomaly-detection-job/main.py", line 49, in main
        kafka_producer = KafkaSink.builder() \
      File "/home/.../.conda/envs/reddit-anomaly-detection-job/lib/python3.9/site-packages/pyflink/datastream/connectors/kafka.py", line 963, in set_record_serializer
        get_field_value(j_topic_selector, 'topicSelector').getClass().getCanonicalName()
    AttributeError: 'NoneType' object has no attribute 'startswith'

涉事代码

# 创建Kafka生产者,使用SimpleStringSchema序列化
record_serializer = KafkaRecordSerializationSchema.builder() \
    .set_topic(kafka_sink_topic) \
    .set_value_serialization_schema(SimpleStringSchema()) \
    .build()

kafka_producer = KafkaSink.builder() \
    .set_bootstrap_servers(bootstrap_servers) \
    .set_record_serializer(record_serializer) \
    .build()

补充信息

相同代码在基于自定义Python镜像的Ververica环境中运行正常,仅本地PyCharm环境出现问题。

解决建议

  • 对齐PyFlink与Flink Java依赖版本:本地环境可能存在版本不匹配问题,PyFlink 1.16必须搭配对应版本的Flink Java组件。执行pip show pyflink确认PyFlink版本,确保本地Flink运行时(或集群)的Java版本与之完全一致。
  • 重建Python虚拟环境:conda环境可能存在依赖冲突,删除现有环境后重新创建:
    conda remove -n reddit-anomaly-detection-job --all
    conda create -n reddit-anomaly-detection-job python=3.9
    conda activate reddit-anomaly-detection-job
    pip install apache-flink==1.16.0
    
  • 更换兼容JDK版本:Flink 1.16官方推荐使用JDK 11,检查本地JDK版本,确保JAVA_HOME环境变量指向正确的JDK路径。
  • 显式添加Kafka连接器依赖:在代码中指定Flink Kafka连接器的JAR包路径,避免本地环境依赖缺失:
    from pyflink.datastream import StreamExecutionEnvironment
    
    env = StreamExecutionEnvironment.get_execution_environment()
    # 替换为本地实际JAR包路径
    env.add_jars(
        "file:///path/to/flink-connector-kafka-1.16.0.jar",
        "file:///path/to/kafka-clients-2.8.1.jar"
    )
    

内容的提问来源于stack exchange,提问作者Monika X

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 11:15:33