Kafka转Spark流处理:DataFrame格式异常问题求助
问题:Kafka转Spark流处理时DataFrame格式异常
在项目中通过Kafka将CSV数据流发送到本地Spark处理,但无法正确解析和展示DataFrame。以下是相关代码:
Kafka生产者代码
import pandas as pd import time import json from confluent_kafka import Producer # Read CSV df = pd.read_csv('timeseries.csv') df = df.astype(str) # Number of rows to send in each batch batch_size = 100 # Kafka settings kafka_bootstrap_servers = 'localhost:9092' kafka_topic = 'T7' # Create a Kafka producer producer_config = {'bootstrap.servers': kafka_bootstrap_servers} producer = Producer(producer_config) # Total number of rows in the dataset total_rows = len(df) # Loop through the dataset and send 100 rows at a time for i in range(0, total_rows, batch_size): batch = df.iloc[i:i + batch_size] # Select only the "id" and "value" columns selected_columns = batch[['id', 'date', 'value', 'label']] # Convert the selected columns to a JSON string json_data = selected_columns.to_json(orient='records') try: # Send the batch to the Kafka topic producer.produce(kafka_topic, key=str(batch['id'].iloc[0]), value=json_data.encode('utf-8')) except Exception as e: print(f"Error sending batch: {e}") # Introduce a delay of 1 second between sending each batch time.sleep(1) # Wait for any outstanding messages to be delivered and delivery reports received producer.flush()
Spark消费者代码(修复后)
from pyspark.sql import SparkSession from pyspark.sql.functions import col, from_json, explode from pyspark.sql.types import StringType, StructType, ArrayType # Define your Spark session spark = SparkSession.builder \ .appName("SparkConsumeKafka") \ .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.0") \ .getOrCreate() # Set log level to ERROR spark.sparkContext.setLogLevel("ERROR") # Read from the Kafka topic - 修正topic与生产者一致 df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "localhost:9092") \ .option("subscribe", "T7") \ .load() # 定义单条记录的Schema my_schema = ( StructType() .add("id", StringType()) .add("date", StringType()) .add("value", StringType()) .add("label", StringType()) ) # 定义数组类型Schema,匹配生产者发送的JSON数组结构 array_schema = ArrayType(my_schema) # 修正解析逻辑:直接解析JSON数组,拆分单条记录 df = df.select( from_json(col("value").cast("string"), array_schema).alias("data_array"), "timestamp" ) df = df.select( explode(col("data_array")).alias("datas"), "timestamp" ) def show_batch(df, epoch_id): df = df.select( "datas.id", "datas.date", "datas.value", "datas.label", "timestamp" ) df.show(truncate=False) # Write to the console with formatted options query = df.writeStream.foreachBatch(show_batch).start() # Await termination query.awaitTermination() # Stop the Spark session spark.stop()
问题原因及修复说明
- Topic不匹配:生产者发送数据到
T7,原消费者订阅T8,导致无法接收正确数据,已修正消费者订阅的topic。 - JSON解析逻辑错误:生产者通过
to_json(orient='records')生成的是JSON数组格式(如[{"id":"1",...}, {"id":"2",...}]),原代码用split("\n")分割的方式不适用,改为先解析为数组类型,再用explode拆分为单条记录。 - Schema补充:新增
ArrayType(my_schema)匹配生产者发送的数组结构,确保解析逻辑与数据格式对齐。
内容的提问来源于stack exchange,提问作者AFORS
相关产品推荐
相关产品推荐

