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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 21:23:17