网易首页 > 网易号 > 正文 申请入驻

使用 ZeroMQ 消息库在 C 和 Python 间共享数据 | Linux 中国

0
分享至

导读:ZeroMQ 是一个快速灵活的消息库,用于数据收集和不同编程语言间的数据共享。

本文字数:12873,阅读时长大约: 15分钟

https://linux.cn/article-12499-1.html
作者:Cristiano L. Fontana
译者:SilentDawn

作为软件工程师,我有多次在要求完成指定任务时感到浑身一冷的经历。其中有一次,我必须在一些新的硬件基础设施和云基础设施之间写一个接口,这些硬件需要 C 语言,而云基础设施主要是用 Python。

实现的方式之一是 ,Python 支持 C 扩展的调用。快速浏览文档后发现,这需要编写大量的 C 代码。这样做的话,在有些情况下效果还不错,但不是我喜欢的方式。另一种方式就是将两个任务放在不同的进程中,并使用 在两者之间交换消息。

在发现 ZeroMQ 之前,遇到这种类型的情况时,我选择了编写扩展的方式。这种方式不算太差,但非常费时费力。如今,为了避免那些问题,我将一个系统细分为独立的进程,通过 发送消息来交换信息。这样,不同的编程语言可以共存,每个进程也变简单了,同时也容易调试。

ZeroMQ 提供了一个更简单的过程:

1. 编写一小段 C 代码,从硬件读取数据,然后把发现的东西作为消息发送出去。

2. 使用 Python 编写接口,实现新旧基础设施之间的对接。

是 ZeroMQ 项目发起者之一,他是个拥有 的非凡人物。

准备

本教程中,需要:

一个 C 编译器(例如 或 )

Fedora 系统上的安装方法:

  1. $ dnf install clang zeromq zeromq-devel python3 python3-zmq

Debian 和 Ubuntu 系统上的安装方法:

  1. $ apt-get install clang libzmq5 libzmq3-dev python3 python3-zmq

如果有问题,参考对应项目的安装指南(上面附有链接)。

编写硬件接口库

因为这里针对的是个设想的场景,本教程虚构了包含两个函数的操作库:

fancyhw_init()用来初始化(设想的)硬件

fancyhw_read_val()用于返回从硬件读取的数据

将库的完整代码保存到文件libfancyhw.h中:

  1. #ifndef LIBFANCYHW_H

  2. #define LIBFANCYHW_H

  3. #include

  4. #include

  5. // This is the fictitious hardware interfacing library

  6. void fancyhw_init(unsigned int init_param)

  7. {

  8. srand(init_param);

  9. }

  10. int16_t fancyhw_read_val(void)

  11. {

  12. return (int16_t)rand();

  13. }

  14. #endif

这个库可以模拟你要在不同语言实现的组件间交换的数据,中间有个随机数发生器。

设计 C 接口

下面从包含管理数据传输的库开始,逐步实现 C 接口。

需要的库

开始先加载必要的库(每个库的作用见代码注释):

  1. // For printf()

  2. #include

  3. // For EXIT_*

  4. #include

  5. // For memcpy()

  6. #include

  7. // For sleep()

  8. #include

  9. #include

  10. #include "libfancyhw.h"

必要的参数

定义main函数和后续过程中必要的参数:

  1. int main(void)

  2. {

  3. const unsigned int INIT_PARAM = 12345;

  4. const unsigned int REPETITIONS = 10;

  5. const unsigned int PACKET_SIZE = 16;

  6. const char *TOPIC = "fancyhw_data";

  7. ...

初始化

所有的库都需要初始化。虚构的那个只需要一个参数:

  1. fancyhw_init(INIT_PARAM);

ZeroMQ 库需要实打实的初始化。首先,定义对象context,它是用来管理全部的套接字的:

  1. void *context = zmq_ctx_new();

  2. if (!context)

  3. {

  4. printf("ERROR: ZeroMQ error occurred during zmq_ctx_new(): %s\n", zmq_strerror(errno));

  5. return EXIT_FAILURE;

  6. }

之后定义用来发送数据的套接字。ZeroMQ 支持若干种套接字,各有其用。使用publish套接字(也叫PUB套接字),可以复制消息并分发到多个接收端。这使得你可以让多个接收端接收同一个消息。没有接收者的消息将被丢弃(即不会入消息队列)。用法如下:

  1. void *data_socket = zmq_socket(context, ZMQ_PUB);

套接字需要绑定到一个具体的地址,这样客户端就知道要连接哪里了。本例中,使用了 (当然也有 ,但 TCP 是不错的默认选择):

  1. const int rb = zmq_bind(data_socket, "tcp://*:5555");

  2. if (rb != 0)

  3. {

  4. printf("ERROR: ZeroMQ error occurred during zmq_ctx_new(): %s\n", zmq_strerror(errno));

  5. return EXIT_FAILURE;

  6. }

下一步, 计算一些后续要用到的值。 注意下面代码中的TOPIC,因为PUB套接字发送的消息需要绑定一个主题。主题用于供接收者过滤消息:

  1. const size_t topic_size = strlen(TOPIC);

  2. const size_t envelope_size = topic_size + 1 + PACKET_SIZE * sizeof(int16_t);

  3. printf("Topic: %s; topic size: %zu; Envelope size: %zu\n", TOPIC, topic_size, envelope_size);

发送消息

启动一个发送消息的循环,循环REPETITIONS次:

  1. for (unsigned int i = 0; i < REPETITIONS; i++)

  2. {

  3. ...

发送消息前,先填充一个长度为PACKET_SIZE的缓冲区。本库提供的是 16 个位的有符号整数。因为 C 语言中int类型占用空间大小与平台相关,不是确定的值,所以要使用指定宽度的int变量:

  1. int16_t buffer[PACKET_SIZE];

  2. for (unsigned int j = 0; j < PACKET_SIZE; j++)

  3. {

  4. buffer[j] = fancyhw_read_val();

  5. }

  6. printf("Read %u data values\n", PACKET_SIZE);

消息的准备和发送的第一步是创建 ZeroMQ 消息,为消息分配必要的内存空间。空白的消息是用于封装要发送的数据的:

  1. zmq_msg_t envelope;

  2. const int rmi = zmq_msg_init_size(&envelope, envelope_size);

  3. if (rmi != 0)

  4. {

  5. printf("ERROR: ZeroMQ error occurred during zmq_msg_init_size(): %s\n", zmq_strerror(errno));

  6. zmq_msg_close(&envelope);

  7. break;

  8. }

现在内存空间已分配,数据保存在 ZeroMQ 消息 “信封”中。函数zmq_msg_data()返回一个指向封装数据缓存区顶端的指针。第一部分是主题,之后是一个空格,最后是二进制数。主题和二进制数据之间的分隔符采用空格字符。需要遍历缓存区的话,使用类型转换和 。(感谢 C 语言,让事情变得直截了当。)做法如下:

  1. memcpy(zmq_msg_data(&envelope), TOPIC, topic_size);

  2. memcpy((void*)((char*)zmq_msg_data(&envelope) + topic_size), " ", 1);

  3. memcpy((void*)((char*)zmq_msg_data(&envelope) + 1 + topic_size), buffer, PACKET_SIZE * sizeof(int16_t))

通过data_socket发送消息:

  1. const size_t rs = zmq_msg_send(&envelope, data_socket, 0);

  2. if (rs != envelope_size)

  3. {

  4. printf("ERROR: ZeroMQ error occurred during zmq_msg_send(): %s\n", zmq_strerror(errno));

  5. zmq_msg_close(&envelope);

  6. break;

  7. }

使用数据之前要先解除封装:

  1. zmq_msg_close(&envelope);

  2. printf("Message sent; i: %u, topic: %s\n", i, TOPIC);

清理

C 语言不提供 功能,用完之后记得要自己扫尾。发送消息之后结束程序之前,需要运行扫尾代码,释放分配的内存:

  1. const int rc = zmq_close(data_socket);

  2. if (rc != 0)

  3. {

  4. printf("ERROR: ZeroMQ error occurred during zmq_close(): %s\n", zmq_strerror(errno));

  5. return EXIT_FAILURE;

  6. }

  7. const int rd = zmq_ctx_destroy(context);

  8. if (rd != 0)

  9. {

  10. printf("Error occurred during zmq_ctx_destroy(): %s\n", zmq_strerror(errno));

  11. return EXIT_FAILURE;

  12. }

  13. return EXIT_SUCCESS;

完整 C 代码

保存下面完整的接口代码到本地名为hw_interface.c的文件:

  1. // For printf()

  2. #include

  3. // For EXIT_*

  4. #include

  5. // For memcpy()

  6. #include

  7. // For sleep()

  8. #include

  9. #include

  10. #include "libfancyhw.h"

  11. int main(void)

  12. {

  13. const unsigned int INIT_PARAM = 12345;

  14. const unsigned int REPETITIONS = 10;

  15. const unsigned int PACKET_SIZE = 16;

  16. const char *TOPIC = "fancyhw_data";

  17. fancyhw_init(INIT_PARAM);

  18. void *context = zmq_ctx_new();

  19. if (!context)

  20. {

  21. printf("ERROR: ZeroMQ error occurred during zmq_ctx_new(): %s\n", zmq_strerror(errno));

  22. return EXIT_FAILURE;

  23. }

  24. void *data_socket = zmq_socket(context, ZMQ_PUB);

  25. const int rb = zmq_bind(data_socket, "tcp://*:5555");

  26. if (rb != 0)

  27. {

  28. printf("ERROR: ZeroMQ error occurred during zmq_ctx_new(): %s\n", zmq_strerror(errno));

  29. return EXIT_FAILURE;

  30. }

  31. const size_t topic_size = strlen(TOPIC);

  32. const size_t envelope_size = topic_size + 1 + PACKET_SIZE * sizeof(int16_t);

  33. printf("Topic: %s; topic size: %zu; Envelope size: %zu\n", TOPIC, topic_size, envelope_size);

  34. for (unsigned int i = 0; i < REPETITIONS; i++)

  35. {

  36. int16_t buffer[PACKET_SIZE];

  37. for (unsigned int j = 0; j < PACKET_SIZE; j++)

  38. {

  39. buffer[j] = fancyhw_read_val();

  40. }

  41. printf("Read %u data values\n", PACKET_SIZE);

  42. zmq_msg_t envelope;

  43. const int rmi = zmq_msg_init_size(&envelope, envelope_size);

  44. if (rmi != 0)

  45. {

  46. printf("ERROR: ZeroMQ error occurred during zmq_msg_init_size(): %s\n", zmq_strerror(errno));

  47. zmq_msg_close(&envelope);

  48. break;

  49. }

  50. memcpy(zmq_msg_data(&envelope), TOPIC, topic_size);

  51. memcpy((void*)((char*)zmq_msg_data(&envelope) + topic_size), " ", 1);

  52. memcpy((void*)((char*)zmq_msg_data(&envelope) + 1 + topic_size), buffer, PACKET_SIZE * sizeof(int16_t));

  53. const size_t rs = zmq_msg_send(&envelope, data_socket, 0);

  54. if (rs != envelope_size)

  55. {

  56. printf("ERROR: ZeroMQ error occurred during zmq_msg_send(): %s\n", zmq_strerror(errno));

  57. zmq_msg_close(&envelope);

  58. break;

  59. }

  60. zmq_msg_close(&envelope);

  61. printf("Message sent; i: %u, topic: %s\n", i, TOPIC);

  62. sleep(1);

  63. }

  64. const int rc = zmq_close(data_socket);

  65. if (rc != 0)

  66. {

  67. printf("ERROR: ZeroMQ error occurred during zmq_close(): %s\n", zmq_strerror(errno));

  68. return EXIT_FAILURE;

  69. }

  70. const int rd = zmq_ctx_destroy(context);

  71. if (rd != 0)

  72. {

  73. printf("Error occurred during zmq_ctx_destroy(): %s\n", zmq_strerror(errno));

  74. return EXIT_FAILURE;

  75. }

  76. return EXIT_SUCCESS;

  77. }

用如下命令编译:

  1. $ clang -std=c99 -I. hw_interface.c -lzmq -o hw_interface

如果没有编译错误,你就可以运行这个接口了。贴心的是,ZeroMQPUB套接字可以在没有任何应用发送或接受数据的状态下运行,这简化了使用复杂度,因为这样不限制进程启动的次序。

运行该接口:

  1. $ ./hw_interface

  2. Topic: fancyhw_data; topic size: 12; Envelope size: 45

  3. Read 16 data values

  4. Message sent; i: 0, topic: fancyhw_data

  5. Read 16 data values

  6. Message sent; i: 1, topic: fancyhw_data

  7. Read 16 data values

  8. ...

  9. ...

输出显示数据已经通过 ZeroMQ 完成发送,现在要做的是让一个程序去读数据。

编写 Python 数据处理器

现在已经准备好从 C 程序向 Python 应用传送数据了。

需要两个库帮助实现数据传输。首先是 ZeroMQ 的 Python 封装:

  1. $ python3 -m pip install zmq

另一个就是 ,用于解码二进制数据。这个库是 Python 标准库的一部分,所以不需要使用pip命令安装。

Python 程序的第一部分是导入这些库:

  1. import zmq

  2. import struct

重要参数

使用 ZeroMQ 时,只能向常量TOPIC定义相同的接收端发送消息:

  1. topic = "fancyhw_data".encode('ascii')

  2. print("Reading messages with topic: {}".format(topic))

初始化

下一步,初始化上下文和套接字。使用subscribe套接字(也称为SUB套接字),它是PUB套接字的天生伴侣。这个套接字发送时也需要匹配主题。

  1. with zmq.Context() as context:

  2. socket = context.socket(zmq.SUB)

  3. socket.connect("tcp://127.0.0.1:5555")

  4. socket.setsockopt(zmq.SUBSCRIBE, topic)

  5. i = 0

  6. ...

接收消息

启动一个无限循环,等待接收发送到SUB套接字的新消息。这个循环会在你按下Ctrl+C组合键或者内部发生错误时终止:

  1. try:

  2. while True:

  3. ... # we will fill this in next

  4. except KeyboardInterrupt:

  5. socket.close()

  6. except Exception as error:

  7. print("ERROR: {}".format(error))

  8. socket.close()

这个循环等待recv()方法获取的新消息,然后将接收到的内容从第一个空格字符处分割开,从而得到主题:

  1. binary_topic, data_buffer = socket.recv().split(b' ', 1)

解码消息

Python 此时尚不知道主题是个字符串,使用标准 ASCII 编解码器进行解码:

  1. topic = binary_topic.decode(encoding = 'ascii')

  2. print("Message {:d}:".format(i))

  3. print("\ttopic: '{}'".format(topic))

下一步就是使用struct库读取二进制数据,它可以将二进制数据段转换为明确的数值。首先,计算数据包中数值的组数。本例中使用的 16 个位的有符号整数对应的是struct中的h

  1. packet_size = len(data_buffer) // struct.calcsize("h")

  2. print("\tpacket size: {:d}".format(packet_size))

知道数据包中有多少组数据后,就可以通过构建一个包含数据组数和数据类型的字符串,来定义格式了(比如“16h”):

  1. struct_format = "{:d}h".format(packet_size)

将二进制数据串转换为可直接打印的一系列数字:

  1. data = struct.unpack(struct_format, data_buffer)

  2. print("\tdata: {}".format(data))

完整 Python 代码

下面是 Python 实现的完整的接收端:

  1. #! /usr/bin/env python3

  2. import zmq

  3. import struct

  4. topic = "fancyhw_data".encode('ascii')

  5. print("Reading messages with topic: {}".format(topic))

  6. with zmq.Context() as context:

  7. socket = context.socket(zmq.SUB)

  8. socket.connect("tcp://127.0.0.1:5555")

  9. socket.setsockopt(zmq.SUBSCRIBE, topic)

  10. i = 0

  11. try:

  12. while True:

  13. binary_topic, data_buffer = socket.recv().split(b' ', 1)

  14. topic = binary_topic.decode(encoding = 'ascii')

  15. print("Message {:d}:".format(i))

  16. print("\ttopic: '{}'".format(topic))

  17. packet_size = len(data_buffer) // struct.calcsize("h")

  18. print("\tpacket size: {:d}".format(packet_size))

  19. struct_format = "{:d}h".format(packet_size)

  20. data = struct.unpack(struct_format, data_buffer)

  21. print("\tdata: {}".format(data))

  22. i += 1

  23. except KeyboardInterrupt:

  24. socket.close()

  25. except Exception as error:

  26. print("ERROR: {}".format(error))

  27. socket.close()

将上面的内容保存到名为online_analysis.py的文件。Python 代码不需要编译,你可以直接运行它。

运行输出如下:

  1. $ ./online_analysis.py

  2. Reading messages with topic: b'fancyhw_data'

  3. Message 0:

  4. topic: 'fancyhw_data'

  5. packet size: 16

  6. data: (20946, -23616, 9865, 31416, -15911, -10845, -5332, 25662, 10955, -32501, -18717, -24490, -16511, -28861, 24205, 26568)

  7. Message 1:

  8. topic: 'fancyhw_data'

  9. packet size: 16

  10. data: (12505, 31355, 14083, -19654, -9141, 14532, -25591, 31203, 10428, -25564, -732, -7979, 9529, -27982, 29610, 30475)

  11. ...

  12. ...

小结

本教程介绍了一种新方式,实现从基于 C 的硬件接口收集数据,并分发到基于 Python 的基础设施的功能。借此可以获取数据供后续分析,或者转送到任意数量的接收端去。它采用了一个消息库实现数据在发送者和处理者之间的传送,来取代同样功能规模庞大的软件。

本教程还引出了我称之为“软件粒度”的概念,换言之,就是将软件细分为更小的部分。这种做法的优点之一就是,使得同时采用不同的编程语言实现最简接口作为不同部分之间沟通的组件成为可能。

实践中,这种设计使得软件工程师能以更独立、合作更高效的方式做事。不同的团队可以专注于数据分析的不同方面,可以选择自己中意的实现工具。这种做法的另一个优点是实现了零代价的并行,因为所有的进程都可以并行运行。 是个令人赞叹的软件,使用它可以让工作大大简化。

via:

作者: 选题: 译者: 校对:

本文由 原创编译, 荣誉推出

特别声明:以上内容(如有图片或视频亦包括在内)为自媒体平台“网易号”用户上传并发布,本平台仅提供信息存储服务。

Notice: The content above (including the pictures and videos if any) is uploaded and posted by a user of NetEase Hao, which is a social media platform and only provides information storage services.

相关推荐
热点推荐
伽马射线异常信号指向暗物质 科学家发现关键疑点

伽马射线异常信号指向暗物质 科学家发现关键疑点

地球观察日记
2026-08-21 02:01:03
村长救助灾民被人举报后被停职处理,洪水再次来袭,村长:管不了

村长救助灾民被人举报后被停职处理,洪水再次来袭,村长:管不了

红豆讲堂
2025-06-21 16:30:05
未来不出意外,那些二胎、三胎的非编工薪家庭,会越来越难,无力托举子女

未来不出意外,那些二胎、三胎的非编工薪家庭,会越来越难,无力托举子女

舒山有鹿
2026-08-15 09:28:35
已有2名伤员出院!升学宴墙体垮塌致5死17伤,主家亲属发声

已有2名伤员出院!升学宴墙体垮塌致5死17伤,主家亲属发声

上观新闻
2026-08-20 06:38:17
不再雪藏败将!国乒亚运名单大调整,王励勤用人引全网两极争议

不再雪藏败将!国乒亚运名单大调整,王励勤用人引全网两极争议

林子说事
2026-08-21 11:00:09
40万亿美元!美国国债爆表,上半年利息就还了5290亿美元,比军费还高!平均每名美国人背负约11.6万美元债务 美国人平均背债11.6万美元

40万亿美元!美国国债爆表,上半年利息就还了5290亿美元,比军费还高!平均每名美国人背负约11.6万美元债务 美国人平均背债11.6万美元

每日经济新闻
2026-08-21 23:06:57
表面上国泰民安,其实暗流涌动!揭秘本轮扫黑升级真正的原因

表面上国泰民安,其实暗流涌动!揭秘本轮扫黑升级真正的原因

王二哥老搞笑
2026-08-11 06:25:10
半年报出炉!长电科技、中国巨石、湖南白银、飞龙股份,后市预期

半年报出炉!长电科技、中国巨石、湖南白银、飞龙股份,后市预期

长风价值掘金
2026-08-21 15:19:53
复刻 B 费神迹!曼联锁定世界级杀器,1.1 亿神锋才是争冠答案

复刻 B 费神迹!曼联锁定世界级杀器,1.1 亿神锋才是争冠答案

一隅非生
2026-08-21 08:43:03
切尔西官方:矿工新赛季欧冠主场将设在斯坦福桥

切尔西官方:矿工新赛季欧冠主场将设在斯坦福桥

懂球帝
2026-08-21 23:11:27
金与正:日本若企图军事扩张,将立即遭到朝鲜毁灭性打击

金与正:日本若企图军事扩张,将立即遭到朝鲜毁灭性打击

鲁中晨报
2026-08-20 09:38:05
张本智和疑似恋情曝光,看台"神秘女孩"刷屏,身份引全网热议

张本智和疑似恋情曝光,看台"神秘女孩"刷屏,身份引全网热议

好乒乓
2026-08-21 21:20:56
善恶有报,移居英国仅2年,57岁吴秀波再迎噩耗,步入李易峰后尘

善恶有报,移居英国仅2年,57岁吴秀波再迎噩耗,步入李易峰后尘

有范又有料
2025-12-17 14:54:06
德国乒乓球权威杂志盛赞樊振东巨星效应,邱党:他才是真正的领袖

德国乒乓球权威杂志盛赞樊振东巨星效应,邱党:他才是真正的领袖

杨华评论
2026-08-21 22:05:01
太心酸!30岁男子失业网贷缠身,硬饿20天险丢命,肾衰竭险些失明

太心酸!30岁男子失业网贷缠身,硬饿20天险丢命,肾衰竭险些失明

王二哥老搞笑
2026-08-22 03:14:36
我被公婆打到住院,丈夫冷笑:只是给你教训,别再闹!第二天他家亲戚都来医院看热闹,到家推开门却全都瘫软

我被公婆打到住院,丈夫冷笑:只是给你教训,别再闹!第二天他家亲戚都来医院看热闹,到家推开门却全都瘫软

晓艾故事汇
2026-08-18 09:28:29
万亿,国产存储终极IPO要来了

万亿,国产存储终极IPO要来了

投资家
2026-08-21 21:42:03
路虎揽胜新车正式公布,8月21日,接受预定

路虎揽胜新车正式公布,8月21日,接受预定

3C毒物
2026-08-22 01:03:47
“长得不好看,真考不上!”老师劝退视频火了,说得已经很委婉了

“长得不好看,真考不上!”老师劝退视频火了,说得已经很委婉了

泽泽先生
2026-08-20 15:35:06
80后电商老板从千万身家到一无所有,负债560万,妻子狠心离婚,老板哭诉自己的失败人生

80后电商老板从千万身家到一无所有,负债560万,妻子狠心离婚,老板哭诉自己的失败人生

捣蛋窝
2026-08-20 09:50:29
2026-08-22 04:35:00
Linux
Linux
Linux 中国开源社区
8018文章数 73111关注度
往期回顾 全部

科技要闻

阿里AI的两面:云赚56亿,千问烧掉138亿

头条要闻

30岁女干部突遇洪水不幸牺牲 诸多细节公布

头条要闻

30岁女干部突遇洪水不幸牺牲 诸多细节公布

体育要闻

续哈登+招沃特森=骑士赌了

娱乐要闻

冯绍峰七夕带儿子游玩!次日聚会

财经要闻

蔡昉解读经济:扩大消费需求的政策思考

汽车要闻

2026成都车展:小鹏全阵容亮相 三款新车套件上新

态度原创

健康
旅游
时尚
艺术
军事航空

“矫正神器”真能治好脊柱侧弯吗?

旅游要闻

苍山洱海闻名天下,这座千年古楼,才是大理真正的历史见证人!

全网全文背诵:阿姨你找谁啊?

艺术要闻

警告:画面极度危险!你的瞳孔即将被这场视觉核爆炸成烟花!

军事要闻

俄在争议岛屿附近进行导弹训练 日本:无法接受

无障碍浏览 进入关怀版