PySpark无法接收Binance Socket端口数据,如何获取最高涨跌幅?
问题解决方案
一、修复Socket数据发送脚本
你的第一个脚本存在两个问题导致PySpark无法正确接收数据:
- 缺少
on_close函数定义:WebSocketApp初始化时指定了on_close但未实现,会导致运行报错 - 未添加换行符分隔消息: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()
关键说明
- 换行符的作用:Spark的
socketTextStream会以换行符作为消息分隔符,发送端必须确保每条JSON消息末尾添加\n,否则Spark会一直等待消息完成。 - 数据解析容错:添加了异常处理,避免单条无效JSON导致整个批次失败。
- 最高涨跌幅计算:使用
transform操作在每个RDD批次中找出涨跌幅最大的记录,若批次无有效数据则返回默认值。
内容的提问来源于stack exchange,提问作者JMFisac
相关产品推荐
相关产品推荐

