PyFlink 1.16 KafkaSink报AttributeError: 'NoneType'无startswith属性
PyFlink 1.16 KafkaSink 本地环境抛出 AttributeError 异常
问题场景
使用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
相关产品推荐
相关产品推荐

