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

PySpark无法接收Binance Socket端口数据,如何获取最高涨跌幅?

问题解决方案

一、修复Socket数据发送脚本

你的第一个脚本存在两个问题导致PySpark无法正确接收数据:

  1. 缺少on_close函数定义:WebSocketApp初始化时指定了on_close但未实现,会导致运行报错
  2. 未添加换行符分隔消息:Spark的socketTextStream默认按行读取数据,每条消息末尾必须添加换行符,否则会被视为单个未完成的消息

修改后的发送脚本:

import json
import socket
import websocket

def on_open(ws, conn):
    print('opened connection')

def on_close(ws, close_status_code, close_msg):
    print('closed connection')

def on_message(ws, message, connexion):
    message = json.loads(message)
    for data in message:
        symbol = data['s']
        price_change_percent = float(data['P'])
        kline_data = {"symbol": symbol, "price_change_percent": price_change_percent}
        json_data = json.dumps(kline_data)
        print(f"Datos enviados a través de la conexión de socket: {json_data}")
        # 添加换行符,确保Spark能按行读取
        result = connexion.send(json_data.encode() + b'\n')
        print(f"Resultado de la llamada a conn.send(): {result}")

def get_data(con):
    socket_address = "wss://stream.binance.com:9443/ws/!ticker@arr"
    ws = websocket.WebSocketApp(
        socket_address,
        on_open=lambda ws: on_open(ws, con),
        on_close=on_close,
        on_message=lambda ws, msg: on_message(ws, msg, connexion=con)
    )
    print("Connecting to Binance WebSocket...")
    ws.run_forever()

TCP_IP = "localhost"
TCP_PORT = 9009
conn = None
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
s.bind((TCP_IP, TCP_PORT))
s.listen(1)
print("Waiting for TCP connection...")
conn, addr = s.accept()
print("Connected... Starting getting klines.")
get_data(conn)

二、修改PySpark脚本实现数据接收与最高涨跌幅计算

修改后的脚本会先解析JSON数据,然后在每个批次中计算涨跌幅最高的交易对:

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

# 创建Spark上下文
sc = SparkContext("local[*]","binanceTickerStream")
sc.setLogLevel("WARN")

# 创建Streaming上下文,批次间隔5秒
ssc = StreamingContext(sc, 5)

# 从Socket读取数据
tickerStream = ssc.socketTextStream("localhost",9009)

# 解析JSON数据,转换为(symbol, price_change_percent)元组
def parse_ticker(json_str):
    try:
        data = json.loads(json_str)
        return (data["symbol"], data["price_change_percent"])
    except:
        return ("invalid", 0.0)

parsed_stream = tickerStream.map(parse_ticker)

# 过滤无效数据
valid_stream = parsed_stream.filter(lambda x: x[0] != "invalid")

# 计算每个批次中涨跌幅最高的记录
def get_max_percent(rdd):
    if not rdd.isEmpty():
        return rdd.max(key=lambda x: x[1])
    else:
        return ("no_data", 0.0)

max_percent_stream = valid_stream.transform(lambda rdd: sc.parallelize([get_max_percent(rdd)]))

# 打印每个批次的结果
print("每个批次涨跌幅最高的交易对:")
max_percent_stream.pprint()

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

关键说明

  1. 换行符的作用:Spark的socketTextStream会以换行符作为消息分隔符,发送端必须确保每条JSON消息末尾添加\n,否则Spark会一直等待消息完成。
  2. 数据解析容错:添加了异常处理,避免单条无效JSON导致整个批次失败。
  3. 最高涨跌幅计算:使用transform操作在每个RDD批次中找出涨跌幅最大的记录,若批次无有效数据则返回默认值。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 01:03:09