【问题标题】:How to transfer a float array (without serializing/deserializing) from Scala (JeroMQ) to C (ZMQ)?如何将浮点数组(不进行序列化/反序列化)从 Scala(JeroMQ)传输到 C(ZMQ)?
【发布时间】:2016-07-11 05:29:06
【问题描述】:

目前,我正在使用 JSON 库在发送方(JeroMQ)序列化数据,并在接收方(C、ZMQ)反序列化。但是,在解析时,JSON 库开始消耗大量内存,操作系统会终止该进程。所以,我想按原样发送浮点数组,即不使用 JSON。

现有的发件人代码如下(syn0syn1Double 数组)。如果syn0syn1 分别约为 100 MB,则进程在解析接收到的数组时被终止,即下面 sn-p 的最后一行:

import org.zeromq.ZMQ
import com.codahale.jerkson
socket.connect("tcp://localhost:5556")

socket.send(json.JSONObject(Map("syn0"->json.JSONArray(List.fromArray(syn0Global)))).toString())
println("SYN0 Request sent”)
val reply_syn0 = socket.recv(0)
println("Response received after syn0: " + new String(reply_syn0))
logInfo("Sending Syn1 request … , size : " + syn1Global.length )

socket.send(json.JSONObject(Map("syn1"->json.JSONArray(List.fromArray(syn1Global)))).toString())
println("SYN1 Request sent")
val reply_syn1 = socket.recv(0)

socket.send(json.JSONObject(Map("foldComplete"->"Done")).toString())
println("foldComplete sent")
//  Get the reply.
val reply_foldComplete = socket.recv(0)
val processedSynValuesJson = new String(reply_foldComplete)
val processedSynValues_jerkson =   jerkson.Json.parse[Map[String,List[Double]]](processedSynValuesJson)

可以不使用 JSON 传输这些数组吗?

这里我在两个 C 程序之间传输一个浮点数组:

//client.c
int main (void)
{
printf ("Connecting to hello world server…\n");
void *context = zmq_ctx_new ();
void *requester = zmq_socket (context, ZMQ_REQ);
zmq_connect (requester, "tcp://localhost:5555");

int request_nbr;
float send_buffer[10];
float recv_buffer[10];

for(int i = 0; i < 10; i++)
    send_buffer[i] = i;

for (request_nbr = 0; request_nbr != 10; request_nbr++) {
    //char buffer [10];
    printf ("Sending Hello %d…\n", request_nbr);
    zmq_send (requester, send_buffer, 10*sizeof(float), 0);
    zmq_recv (requester, recv_buffer, 10*sizeof(float), 0);
    printf ("Received World %.3f\n", recv_buffer[5]);
}
zmq_close (requester);
zmq_ctx_destroy (context);
return 0;
}

//server.c

int main (void)
{
//  Socket to talk to clients
void *context = zmq_ctx_new ();
void *responder = zmq_socket (context, ZMQ_REP);
int rc = zmq_bind (responder, "tcp://*:5555");
assert (rc == 0);
float recv_buffer[10];
float send_buffer[10];
while (1) {
    //char buffer [10];
    zmq_recv (responder, recv_buffer, 10*sizeof(float), 0);
    printf ("Received Hello\n");
    for(int i = 0; i < 10; i++)
            send_buffer[i] = recv_buffer[i]+5;
    zmq_send (responder, send_buffer, 10*sizeof(float), 0);
}
return 0;
}

最后,我尝试使用 Scala 做类似的事情没有成功(下面是客户端代码):

def main(args: Array[String]) {
val context = ZMQ.context(1)
val socket = context.socket(ZMQ.REQ)

println("Connecting to hello world server…")
socket.connect ("tcp://localhost:5555")
val msg : Array[Float] = Array(1,2,3,4,5,6,7,8,9,10)
val bbuf = java.nio.ByteBuffer.allocate(4*msg.length)
bbuf.asFloatBuffer.put(java.nio.FloatBuffer.wrap(msg))


for (request_nbr <- 1 to 10)  {
    socket.sendByteBuffer(bbuf,0)

}
}

【问题讨论】:

  • 您使用哪种语言?你说 JSON 不起作用,你在哪里尝试以二进制形式发送它们?您可以以任何您喜欢的方式发送数据。
  • 斯卡拉。请查看更新后的帖子。

标签: json scala zeromq distributed jeromq


【解决方案1】:

SER/DES ?尺寸?
不,与运输哲学相关的潜在约束很重要。

您已开始使用 0.1 GB 大小调整传输负载,并报告了 JSON-library 分配导致您的操作系统终止进程。

接下来,在其他帖子中,您请求了 0.762 GB 大小以用于传输负载。

但在ZeroMQ 传输编排中还有一个比选择外部数据序列化器SER/DES 策略更重要的问题。

没有人会禁止您尝试发送尽可能大的 BLOB,而 JSON-decorated 字符串已经向您展示了此类方法的阴暗面,还有其他不继续往前走的理由。

ZeroMQ 毫无疑问是一个伟大而强大的工具箱。仍然需要一些时间才能获得真正智能且高性能的代码部署所必需的洞察力,从而最大限度地利用这个强大的主力。

功能丰富的内部生态系统“幕后”的副作用之一是隐藏在消息传递概念中的鲜为人知的策略。

一个人可以发送任何大小合理的消息,但不能保证送达。它要么完全交付,或什么都没有完全交付,如上所述,没有任何保证。

哎哟?!

是的,不保证。

基于这一核心零保证理念,在决定步骤和措施时应格外小心,如果您打算在 Gigabyte BEASTs 之间来回移动,则更应谨慎。

从这个意义上说,真正的 SUT 测试可能会在数量上支持小型消息可以传输(如果您确实仍然需要移动 GB(请参阅上面的评论,在OP ) 并且别无选择 ) 将整个数据量分割成更小的部分,并采用容易出错的重新组装措施,这导致比端到端解决方案更快、更安全尝试使用哑力并指示代码将大约 1 GB 的数据转储到实际可用的任何资源上(ZeroMQ 的零复制原则不能也不会本身为您节省这些努力)。

有关另一个隐藏陷阱的详细信息,与不完全零复制实现有关,read Martin SUSTRIK's, co-father of ZeroMQ, remarks on Zero-Copy "till-kernel-boundary-only"(因此,至少预期的内存空间分配会增加一倍...)。


解决方案:

重新设计架构以传播小型消息,如果不在远程进程中保持“镜像”原始数据结构,而不是尝试保持一次性千兆传输的可存活性。


最好的下一步?

虽然用几个 SLOC-s 并不能解决您的问题,但最好的办法是,如果您认真地将您的智力投入到分布式处理中,那就是阅读 Pieter HINTJEN 的可爱书“代码连接,第 1 卷”

是的,产生自己的洞察力需要一些时间,但这会在许多方面将您提升到另一个专业代码设计水平。 值得花时间。值得努力。

【讨论】:

    【解决方案2】:

    您需要以某种形式或方式序列化数据 - 最终,您在内存中获取一个结构,并指导另一侧如何重建该结构(使用两种不同语言的奖励点,其中无论如何,内存中的结构可能不同)。我建议您使用新的 JSON 库,因为这似乎是问题所在,但您可以使用更有效的协议。 Protocol Buffers 享受多种语言的良好支持,这可能是我要开始的地方。

    【讨论】:

    • 你能解释一下为什么浮点数数组不能简单地作为字节序列传输吗?我写了两个 C 程序,它们使用 ZMQ 来简单地传输一个浮点数组,它工作正常。现在,如果我用 Scala 客户端替换 C 客户端,为什么它不起作用。请参阅更新帖子中的 C 代码,以及我在 Scala 中编写客户端的尝试。
    • 好的,所以看看您更新的问题,缓冲区是您对数据的临时“序列化”。我的假设是缓冲区的实现在 C 和 Java 之间有所不同——这是完全合理的。字节只是不对齐,C 不知道如何处理数据。
    • 这就是我的困惑所在。你是说即使我通过 ZMQ 发送一个字节块,也有一些附加的元数据,只能由使用相同语言实现的客户端解释?你能指点参考吗?
    • 这不是我要说的,但我正在深入研究 libzmq 和 jeromq 的代码,这样我就可以更权威地谈论正在发生的事情,而不是推测。
    • 没有为 C 和 Scala 设置我自己的环境(我两者都没有开发)我不能给你一个明确的答案,但我的建议是表示由 C 创建的浮点数组的字节流是与表示由 Scala 创建的浮点数组的字节流不同。我可以告诉你,在 C 中不知道数组是 float 数组 - zmq_send() 所做的第一件事是 memcpy()according to this 将其转换为无符号字符数组。最好的办法是尝试打印出字节表示并查看自己
    猜你喜欢
    • 2013-07-17
    • 1970-01-01
    • 2015-02-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-08-10
    相关资源
    最近更新 更多