在ODI 12C中以Spark为中转从Kafka映射到CSV文件时遇错误
解决ODI 12C中Spark中转Kafka到CSV时的AttributeError问题
这个错误我之前帮人排查过,核心是PySpark Streaming Kafka的参数传错了,咱们一步步来理清楚:
错误原因拆解
从报错里的代码片段能看到关键问题:
,"group.id" : "kafka_file_Physical"}, fromOffsets='earliest')
这里的fromOffsets='earliest'完全用错了地方:
- PySpark的
createDirectStream方法要求fromOffsets是字典类型,格式是{主题名: 分区偏移量},用来指定特定分区的起始消费位置; - 你传了字符串
'earliest',而字符串根本没有items()方法(Spark内部会调用这个方法解析偏移量配置),所以直接抛出了AttributeError: 'str' object has no attribute 'items'。
如果你的需求是从Kafka的最早可用偏移量开始消费,这个配置应该放到kafkaParams字典里,而不是传给fromOffsets。
修复方案
方案1:从最早偏移量开始消费(最常见需求)
修改PySpark代码的参数配置,把auto.offset.reset加到kafka参数中,移除错误的fromOffsets='earliest':
# 正确的Kafka参数配置 kafkaParams = { "group.id": "kafka_file_Physical", "auto.offset.reset": "earliest" # 这里配置从最早偏移量启动消费 } # 调用createDirectStream时无需传fromOffsets stream = KafkaUtils.createDirectStream(ssc, ["你的主题名"], kafkaParams)
方案2:指定具体分区的偏移量(如果需要精准控制)
如果你确实需要从某个分区的特定位置开始消费,fromOffsets必须传字典,示例如下:
from pyspark.streaming.kafka import TopicAndPartition # 定义特定主题分区的起始偏移量 fromOffsets = { TopicAndPartition("你的主题名", 0): 100, # 0号分区从第100条开始 TopicAndPartition("你的主题名", 1): 200 # 1号分区从第200条开始 } # 调用时传入正确的字典格式参数 stream = KafkaUtils.createDirectStream(ssc, ["你的主题名"], kafkaParams, fromOffsets=fromOffsets)
最后回到ODI 12C里,找到生成这个PySpark脚本的Spark EKM配置节点,调整Kafka消费的偏移量配置方式,确保参数传递符合PySpark的要求就行。
内容的提问来源于stack exchange,提问作者rounak chourasia
相关产品推荐
相关产品推荐

