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

如何将TCP流式传输的JSON数据转换为Spark DataFrame

问题解决方案

1 修正stream.py数据输出格式

Spark Streaming的socketTextStream默认按行切分数据,传输的每条JSON必须以换行符\n结尾,否则会出现JSON不完整、解析失败的问题。核心发送逻辑示例:

# 单条数据格式参考:{"0": {"feature0": 1, "feature1": "测试文本"}}
import socket
import json
import time

server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
server_socket.bind(("localhost", 6100))
server_socket.listen(1)
client, addr = server_socket.accept()

batch = [
    {"0": {"feature0": 0, "feature1": "这是负面文本"}},
    {"1": {"feature0": 1, "feature1": "这是正面文本"}}
]
for item in batch:
    # 必须加换行符
    client.sendall((json.dumps(item) + "\n").encode("utf-8"))
    time.sleep(0.5)

2 修正Spark Streaming客户端代码

常见问题点包括:Spark master配置线程不足、未过滤空RDD、解析逻辑未处理异常、未提前定义Schema。完整可运行代码:

from pyspark import SparkContext
from pyspark.streaming import StreamingContext
from pyspark.sql import SparkSession
import json

# 配置local[2]及以上,预留1个线程接收数据、1个线程处理数据
sc = SparkContext("local[2]", "JSONStreamProcess")
# 批次间隔设为2秒,不要设置过小导致未收集到数据就触发处理
ssc = StreamingContext(sc, 2)
spark = SparkSession(sc)

# 提前定义Schema,避免自动推断出错
df_schema = "`batch_index` INT, feature0 INT, feature1 STRING"

def process_batch_rdd(rdd):
    # 跳过空批次,避免无意义报错
    if rdd.isEmpty():
        return
    # 解析+过滤异常数据,扁平化外层序号结构
    parsed_rdd = rdd.map(lambda line: json.loads(line.strip())) \
                    .filter(lambda x: isinstance(x, dict) and len(x) == 1) \
                    .map(lambda x: {
                        "batch_index": int(list(x.keys())[0]),
                        **list(x.values())[0]
                    })
    # 转换为DataFrame并输出
    result_df = spark.createDataFrame(parsed_rdd, schema=df_schema)
    result_df.show(truncate=False)

# 连接socket流
stream_lines = ssc.socketTextStream("localhost", 6100)
stream_lines.foreachRDD(process_batch_rdd)

# 启动流处理
ssc.start()
ssc.awaitTermination()

3 运行与排查步骤

  • 先启动stream.py,再启动Spark Streaming客户端,顺序颠倒会导致连接失败无数据
  • 启动前确认6100端口未被占用,可执行对应系统命令验证:
    • Linux/macOS:lsof -i:6100
    • Windows:netstat -ano | findstr 6100
  • 若仍无输出,可在process_batch_rdd中新增打印逻辑验证原始数据接收情况:rdd.foreach(lambda line: print("原始接收内容:", line)),如果能打印原始内容就是JSON解析逻辑问题,不能打印就是socket连接或stream.py发送逻辑问题。

内容的提问来源于stack exchange,提问作者Achyut Jagini

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 01:54:03