使用Kafka、Python和Spark处理XML流数据遇AttributeError问题求助
问题分析与解决
错误根源
你遇到的AttributeError: 'DataFrameWriter' object has no attribute 'start'是因为:
- 你用批处理API(
spark.read)读取Kafka,得到的是静态DataFrame,对应的DataFrameWriter只支持show()、save()等批处理输出方法,start()是流处理API(spark.readStream)专属方法,不能混用。 - 代码里冗余创建了
kafka.KafkaConsumer实例,完全没用,直接删掉即可。 - 你要解析XML数据,但当前代码用
from_json处理,这是错误的——XML和JSON格式完全不同,需要用Spark XML解析库处理。
修正方案
方案1:批处理读取Kafka并解析XML
如果需求是一次性读取Kafka中的现有数据(批处理),修改代码如下:
# 导入必要包 from pyspark.sql import SparkSession from pyspark.sql.types import StructType, IntegerType, StringType from pyspark.sql.functions import col, from_xml # 初始化SparkSession,引入XML解析依赖(版本需匹配你的Spark版本) spark = SparkSession\ .builder\ .master("local[1]")\ .appName("Consumer")\ .config("spark.jars.packages", "com.databricks:spark-xml_2.12:0.16.0") .getOrCreate() topic_Name = 'XML_File_Processing3' # 批处理读取Kafka数据 kafka_df = spark\ .read \ .format("kafka") \ .option("kafka.bootstrap.servers", "localhost:9092") \ .option("failOnDataLoss", "false") \ .option("subscribe", topic_Name) \ .load() # 将二进制value转为字符串 new_df = kafka_df.selectExpr("CAST(value AS STRING)") # 定义XML对应Schema(根据你的XML结构调整) book_schema = StructType()\ .add("book_id", IntegerType())\ .add("author", StringType())\ .add("title", StringType())\ .add("genre", StringType())\ .add("price", IntegerType())\ .add("publish_date", IntegerType())\ .add("description", StringType()) # 解析XML字符串为DataFrame book_DF = new_df.select(from_xml(col("value"), book_schema).alias("dataf")) # 批处理输出到控制台,直接用show() book_DF.show(truncate=False)
方案2:流处理读取Kafka并解析XML
如果需求是实时监听Kafka主题(流处理),修改代码如下:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, IntegerType, StringType from pyspark.sql.functions import col, from_xml # 初始化SparkSession,引入XML解析依赖 spark = SparkSession\ .builder\ .master("local[1]")\ .appName("Consumer")\ .config("spark.jars.packages", "com.databricks:spark-xml_2.12:0.16.0") .getOrCreate() topic_Name = 'XML_File_Processing3' # 流处理方式读取Kafka数据 kafka_stream_df = spark\ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "localhost:9092") \ .option("failOnDataLoss", "false") \ .option("subscribe", topic_Name) \ .load() new_stream_df = kafka_stream_df.selectExpr("CAST(value AS STRING)") # 定义XML Schema book_schema = StructType()\ .add("book_id", IntegerType())\ .add("author", StringType())\ .add("title", StringType())\ .add("genre", StringType())\ .add("price", IntegerType())\ .add("publish_date", IntegerType())\ .add("description", StringType()) book_stream_DF = new_stream_df.select(from_xml(col("value"), book_schema).alias("dataf")) # 流处理输出到控制台,用writeStream和start() query = book_stream_DF.writeStream\ .format("console")\ .outputMode("append")\ .start() query.awaitTermination()
注意事项
spark-xml版本要和你的Spark版本匹配,比如Spark 3.3.x对应spark-xml_2.12:0.16.0。- 若你的XML是根节点包含多个book的数组结构,需将Schema改为
ArrayType(book_schema)并对应调整解析逻辑。 - 若你的Kafka集群未启用SSL,直接删除
kafka.security.protocol配置项。
内容的提问来源于stack exchange,提问作者Monam Bharti
相关产品推荐
相关产品推荐

