【问题标题】:Python - Stream - How to send data to server fasterPython - Stream - 如何更快地将数据发送到服务器
【发布时间】:2018-06-12 23:43:20
【问题描述】:

我想将文本文件中的字符串发送到我的本地端口,但我必须为每个字符串打开连接并关闭它。因此,数据流很慢。(几乎在两秒内1个字符串)。我怎样才能让它更快?

        while 1:
            conn, addr = s.accept()

            line_number = random.randint(1,2261074)
            liste.append(line_number)
            line = linecache.getline(filename,line_number)

            sendit = line.split("   ")[1]
            print type(sendit)
            print "sending: " + sendit
            conn.send(sendit)

            conn.close()
            print('End Of Stream.')

【问题讨论】:

  • 减慢发送速度的主要因素是您总是打开新连接。这在很大程度上受到网络延迟的限制。因此,只有放弃此概念并始终使用相同的 TCP 连接,才有可能实现重大加速。但“我必须为每个字符串打开连接并关闭它” 建议您必须遵守此行为。在这种情况下:不可能有很大的加速。
  • 实际上我的主要目的是将数据发送到我的端口,但火花流无法接收数据。所以,在这里,link 他们建议我每次都打开和关闭连接。
  • 您尝试打开一次,发送所有内容,然后关闭。现在您在发送每个字符串后打开和关闭。有没有什么方法可以在每发送 N 行后打开和关闭连接?

标签: python sockets apache-spark stream spark-streaming


【解决方案1】:

answer 建议在每个连接上向 Spark 发送 10 条消息,将每条消息间隔 1 秒,然后关闭连接。

保持连接打开并在服务器端使用非阻塞套接字可能会更好。

下面的服务器代码保持连接打开,在非阻塞套接字上分批发送消息,每批之间有一个空闲延迟。

这可以用来测试 Spark 接收消息的速度。我已将其设置为分批发送 50 条消息,然后等待 1 秒再发送下一个 50 条。

即使我将空闲延迟设置为零,Spark 也可以在我的机器上接收到所有消息。

您可以根据应用的需要进行试验和调整。

服务器代码:

import socket
import time
import select

def do_send(sock, msg, timeout):
    readers = []
    writers = [sock]
    excepts = [sock]
    rxs, txs, exs = select.select(readers, writers, excepts, timeout)
    if sock in exs:
        return False
    elif sock in txs:
        sock.send(msg)
        return True
    else:
        return False

host = 'localhost'
port = 9999

s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
s.bind((host, port))
s.listen(1)

print "waiting for client"
conn, addr = s.accept()
print "client connected"
conn.setblocking(0)

batchSize = 50
idle = 1.
count = 0
running = 1
while running:
    try:
        sc = 0
        while (sc < batchSize):
            if do_send (conn, 'Hello World ' + repr(count) + '\n', 1.0):
                sc += 1
                count += 1
        print "sent " + repr(batchSize) + ", waiting " + repr(idle) + " seconds"
        time.sleep(idle)
    except socket.error:
        conn.close()
        running = 0

print "done"

简单的 Spark 代码:

from pyspark import SparkContext
from pyspark.streaming import StreamingContext

sc = SparkContext("local[2]","test")
ssc = StreamingContext(sc, 1)
ststr = ssc.socketTextStream("localhost", 9999)
lines = ststr.flatMap(lambda line: line.split('\n'))
lines.pprint()
ssc.start()  
ssc.awaitTermination()

希望这可能有用。

【讨论】:

  • 不,这个也不起作用。首先,它看起来不像流媒体。当我打印流数据时。它从第一个数据打印到最后一个数据。但实际上 spark 必须只打印在特定间隔时间内到达的数据。其次,地图在该行不起作用,lines = ststr.flatMap(lambda line: line.split('\n')).map(lambda x: (x,1))
  • @osman tamer 很好,它对我有用(但我知道这对你没有帮助)。如果您需要更多帮助,您需要发布完整的代码以供他人查看,而不仅仅是 while 1:sn-p,并解释您正在寻求的结果。
猜你喜欢
  • 2014-11-17
  • 2019-05-27
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-04-14
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多