- ubuntu12.04环境下使用kvm ioctl接口实现最简单的虚拟机
- Ubuntu 通过无线网络安装Ubuntu Server启动系统后连接无线网络的方法
- 在Ubuntu上搭建网桥的方法
- ubuntu 虚拟机上网方式及相关配置详解
CFSDN坚持开源创造价值,我们致力于搭建一个资源共享平台,让每一个IT人在这里找到属于你的精彩世界.
这篇CFSDN的博客文章Python的消息队列包SnakeMQ使用初探由作者收集整理,如果你对这篇文章有兴趣,记得点赞哟.
1、关于snakemq的官方介绍 SnakeMQ的GitHub项目页:https://github.com/dsiroky/snakemq 1.纯python实现,跨平台 。
2.自动重连接 。
3.可靠发送--可配置的消息方式与消息超时方式 。
4.持久化/临时 两种队列 。
5.支持异步 -- poll() 。
6.symmetrical -- 单个TCP连接可用于双工通讯 。
7.多数据库支持 -- SQLite、MongoDB…… 。
8.brokerless - 类似ZeroMQ的实现原理 。
9.扩展模块:RPC, bandwidth throttling 。
以上都是官话,需要自己验证,动手封装了一下,感觉萌萌哒.
。
2、几个主要问题说明 。
1.支持自动重连,不需要自己动手写心跳逻辑,你只需要关注发送和接收就行 。
2.支持数据持久化,如果开始持久化,在重连之后会自动发送数据.
3.数据的接收,snakemq通过提供回调实现,你只需要写个接收方法添加到回调列表里去.
4.数据的发送,在此发送的都是bytes类型(二进制),因此需要转换。我在程序中测试的都是文本字符串,使用str.encode(‘utf-8')转换成bytes,接收时再转换回来.
5.术语解释,Connector:类似于socket的TcpClient,Lisenter:类似于socket的TcpServer,每个connector或者listener都一个一个ident标识,发送和接收数据时就知道是谁的数据了.
6.使用sqlite持久化时,需要修改源码,sqlite3.connect(filename,check_same_thread = False),用于解决多线程访问sqlite的问题。(会不会死锁?) 。
7.启动持久化时,如果重新连上,则会自动发送,保证可靠.
8.为了封装的需要,数据接收以后,我通过callback方式传送出去.
。
3、代码 。
说明代码中使用了自定义的日志模块 。
1
2
3
|
from
common
import
nxlogger
import
snakemqlogger as logger
|
可替换成logging的.
回调类(callbacks.py):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
|
# -*- coding:utf-8 -*-
'''synchronized callback'''
class
Callback(
object
):
def
__init__(
self
):
self
.callbacks
=
[]
def
add(
self
, func):
self
.callbacks.append(func)
def
remove(
self
, func):
self
.callbacks.remove(func)
def
__call__(
self
,
*
args,
*
*
kwargs):
for
callback
in
self
.callbacks:
callback(
*
args,
*
*
kwargs)
|
Connector类(snakemqConnector.py):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
|
# -*- coding:utf-8 -*-
import
threading
import
snakemq
import
snakemq.link
import
snakemq.packeter
import
snakemq.messaging
import
snakemq.message
from
snakemq.storage.sqlite
import
SqliteQueuesStorage
from
snakemq.message
import
FLAG_PERSISTENT
from
common.callbacks
import
Callback
from
common
import
nxlogger
import
snakemqlogger as logger
class
SnakemqConnector(threading.Thread):
def
__init__(
self
, snakemqident
=
None
, remoteIp
=
"localhost"
, remotePort
=
9090
, persistent
=
False
):
super
(SnakemqConnector,
self
).__init__()
self
.messaging
=
None
self
.link
=
None
self
.snakemqident
=
snakemqident
self
.pktr
=
None
self
.remoteIp
=
remoteIp
self
.remotePort
=
remotePort
self
.persistent
=
persistent
self
.on_recv
=
Callback()
self
._initConnector()
def
run(
self
):
logger.info(
"connector start..."
)
if
self
.link !
=
None
:
self
.link.loop()
logger.info(
"connector end..."
)
def
terminate(
self
):
logger.info(
"connetor terminating..."
)
if
self
.link !
=
None
:
self
.link.stop()
self
.link.cleanup()
logger.info(
"connetor terminated"
)
def
on_recv_message(
self
, conn, ident, message):
try
:
self
.on_recv(ident, message.data.decode(
'utf-8'
))
#dispatch received data
except
Exception as e:
logger.error(
"connector recv:{0}"
.
format
(e))
print
(e)
'''send message to dest host named destIdent'''
def
sendMsg(
self
, destIdent, byteseq):
msg
=
None
if
self
.persistent:
msg
=
snakemq.message.Message(byteseq, ttl
=
60
, flags
=
FLAG_PERSISTENT)
else
:
msg
=
snakemq.message.Message(byteseq, ttl
=
60
)
if
self
.messaging
=
=
None
:
logger.error(
"connector:messaging is not initialized, send message failed"
)
return
self
.messaging.send_message(destIdent, msg)
'''
'''
def
_initConnector(
self
):
try
:
self
.link
=
snakemq.link.Link()
self
.link.add_connector((
self
.remoteIp,
self
.remotePort))
self
.pktr
=
snakemq.packeter.Packeter(
self
.link)
if
self
.persistent:
storage
=
SqliteQueuesStorage(
"SnakemqStorage.db"
)
self
.messaging
=
snakemq.messaging.Messaging(
self
.snakemqident, "",
self
.pktr, storage)
else
:
self
.messaging
=
snakemq.messaging.Messaging(
self
.snakemqident, "",
self
.pktr)
self
.messaging.on_message_recv.add(
self
.on_recv_message)
except
Exception as e:
logger.error(
"connector:{0}"
.
format
(e))
finally
:
logger.info(
"connector[{0}] loop ended..."
.
format
(
self
.snakemqident))
|
Listener类(snakemqListener.py):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
|
# -*- coding:utf-8 -*-
import
threading
import
snakemq
import
snakemq.link
import
snakemq.packeter
import
snakemq.messaging
import
snakemq.message
from
common
import
nxlogger
import
snakemqlogger as logger
from
common.callbacks
import
Callback
class
SnakemqListener(threading.Thread):
def
__init__(
self
, snakemqident
=
None
, ip
=
"localhost"
, port
=
9090
, persistent
=
False
):
super
(SnakemqListener,
self
).__init__()
self
.messaging
=
None
self
.link
=
None
self
.pktr
=
None
self
.snakemqident
=
snakemqident
self
.ip
=
ip;
self
.port
=
port
self
.connectors
=
{}
self
.on_recv
=
Callback()
self
.persistent
=
persistent
self
._initlistener()
'''
thread run
'''
def
run(
self
):
logger.info(
"listener start..."
)
if
self
.link !
=
None
:
self
.link.loop()
logger.info(
"listener end..."
)
'''
terminate snakemq listener thread
'''
def
terminate(
self
):
logger.info(
"listener terminating..."
)
if
self
.link !
=
None
:
self
.link.stop()
self
.link.cleanup()
logger.info(
"listener terminated"
)
'''
receive message from host named ident
'''
def
on_recv_message(
self
, conn, ident, message):
try
:
self
.on_recv(ident, message.data.decode(
'utf-8'
))
#dispatch received data
self
.sendMsg(
'bob'
,
'hello,{0}'
.
format
(ident).encode(
'utf-8'
))
except
Exception as e:
logger.error(
"listener recv:{0}"
.
format
(e))
print
(e)
def
on_drop_message(
self
, ident, message):
print
(
"message dropped"
, ident, message)
logger.debug(
"listener:message dropped,ident:{0},message:{1}"
.
format
(ident, message))
'''client connect'''
def
on_connect(
self
, ident):
logger.debug(
"listener:{0} connected"
.
format
(ident))
self
.connectors[ident]
=
ident
self
.sendMsg(ident,
"hello"
.encode(
'utf-8'
))
'''client disconnect'''
def
on_disconnect(
self
, ident):
logger.debug(
"listener:{0} disconnected"
.
format
(ident))
if
ident
in
self
.connectors:
self
.connectors.pop(ident)
'''
listen start loop
'''
def
_initlistener(
self
):
try
:
self
.link
=
snakemq.link.Link()
self
.link.add_listener((
self
.ip,
self
.port))
self
.pktr
=
snakemq.packeter.Packeter(
self
.link)
self
.pktr.on_connect.add(
self
.on_connect)
self
.pktr.on_disconnect.add(
self
.on_disconnect)
if
self
.persistent:
storage
=
SqliteQueuesStorage(
"SnakemqStorage.db"
)
self
.messaging
=
snakemq.messaging.Messaging(
self
.snakemqident, "",
self
.pktr, storage)
else
:
self
.messaging
=
snakemq.messaging.Messaging(
self
.snakemqident, "",
self
.pktr)
self
.messaging.on_message_recv.add(
self
.on_recv_message)
self
.messaging.on_message_drop.add(
self
.on_drop_message)
except
Exception as e:
logger.error(
"listener:{0}"
.
format
(e))
finally
:
logger.info(
"listener:loop ended..."
)
'''send message to dest host named destIdent'''
def
sendMsg(
self
, destIdent, byteseq):
msg
=
None
if
self
.persistent:
msg
=
snakemq.message.Message(byteseq, ttl
=
60
, flags
=
FLAG_PERSISTENT)
else
:
msg
=
snakemq.message.Message(byteseq, ttl
=
60
)
if
self
.messaging
=
=
None
:
logger.error(
"listener:messaging is not initialized, send message failed"
)
return
self
.messaging.send_message(destIdent, msg)
|
测试代码connector(testSnakeConnector.py):
读取本地一个1M的文件,然后发送给listener,然后listener发回一个hello的信息.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
|
from
netComm.snakemq
import
snakemqConnector
import
time
import
sys
import
os
def
received(ident, data):
print
(data)
if
__name__
=
=
"__main__"
:
bob
=
snakemqConnector.SnakemqConnector(
'bob'
,
"10.16.5.45"
,
4002
,
True
)
bob.on_recv.add(received)
bob.start()
try
:
with
open
(
"testfile.txt"
,encoding
=
'utf-8'
) as f:
txt
=
f.read()
for
i
in
range
(
100
):
bob.sendMsg(
"niess"
,txt.encode(
'utf-8'
))
time.sleep(
0.1
)
except
Exception as e:
print
(e)
time.sleep(
5
)
bob.terminate()
测试代码listener(testSnakeListener.py):
from
netComm.snakemq
import
snakemqListener
import
time
def
received(ident, data):
filename
=
"log/recFile{0}.txt"
.
format
(time.strftime(
'%S'
,time.localtime()))
file
=
open
(filename,
'w'
)
file
.writelines(data)
file
.close()
if
__name__
=
=
"__main__"
:
niess
=
snakemqListener.SnakemqListener(
"niess"
,
"10.16.5.45"
,
4002
)
niess.on_recv.add(received)
niess.start()
print
(
"niess start..."
)
time.sleep(
60
)
niess.terminate()
print
(
"niess end..."
)
|
最后此篇关于Python的消息队列包SnakeMQ使用初探的文章就讲到这里了,如果你想了解更多关于Python的消息队列包SnakeMQ使用初探的内容请搜索CFSDN的文章或继续浏览相关文章,希望大家以后支持我的博客! 。
我一直在读到,如果一个集合“被释放”,它也会释放它的所有对象。另一方面,我还读到,一旦集合被释放,集合就会释放它的对象。 但最后一件事可能并不总是发生,正如苹果所说。系统决定是否取消分配。在大多数情况
我有一个客户端-服务器应用程序,它使用 WCF 进行通信,并使用 NetDataContractSerializer 序列化对象图。 由于服务器和客户端之间传输了大量数据,因此我尝试通过微调数据成员的
我需要有关 JMS 队列和消息处理的帮助。 我有一个场景,需要针对特定属性组同步处理消息,但可以在不同属性组之间同时处理消息。 我了解了特定于每个属性的消息组和队列的一些知识。我的想法是,我想针对
我最近开始使用 C++,并且有一种强烈的冲动 #define print(msg) std::cout void print(T const& msg) { std::cout void
我已经为使用 JGroups 编写了简单的测试。有两个像这样的简单应用程序 import org.jgroups.*; import org.jgroups.conf.ConfiguratorFact
这个问题在这里已经有了答案: Firebase messaging is not supported in your browser how to solve this? (3 个回答) 7 个月前关
在我的 C# 控制台应用程序中,我正在尝试更新 CRM 2016 中的帐户。IsFaulted 不断返回 true。当我向下钻取时它返回的错误消息如下: EntityState must be set
我正在尝试通过 tcp 将以下 json 写入 graylog 服务器: {"facility":"GELF","file":"","full_message":"Test Message Tcp",
我正在使用 Django 的消息框架来指示成功的操作和失败的操作。 如何排除帐户登录和注销消息?目前,登录后登陆页面显示 已成功登录为“用户名”。我不希望显示此消息,但应显示所有其他成功消息。我的尝试
我通过编写禁用qDebug()消息 CONFIG(release, debug|release):DEFINES += QT_NO_DEBUG_OUTPUT 在.pro文件中。这很好。我想知道是否可以
我正在使用 ThrottleRequest 来限制登录尝试。 在 Kendler.php 我有 'throttle' => \Illuminate\Routing\Middleware\Throttl
我有一个脚本,它通过die引发异常。捕获异常时,我想输出不附加位置信息的消息。 该脚本: #! /usr/bin/perl -w use strict; eval { die "My erro
允许的消息类型有哪些(字符串、字节、整数等)? 消息的最大大小是多少? 队列和交换器的最大数量是多少? 最佳答案 理论上任何东西都可以作为消息存储/发送。实际上您不想在队列上存储任何内容。如果队列大部
基本上,我正在尝试创建一个简单的 GUI 来与 Robocopy 一起使用。我正在使用进程打开 Robocopy 并将输出重定向到文本框,如下所示: With MyProcess.StartI
我想将进入 MQ 队列的消息记录到数据库/文件或其他日志队列,并且我无法修改现有代码。是否有任何方法可以实现某种类似于 HTTP 嗅探器的消息记录实用程序?或者也许 MQ 有一些内置的功能来记录消息?
我得到了一个带有 single_selection 数据表和一个命令按钮的页面。命令按钮调用一个 bean 方法来验证是否进行了选择。如果不是,它应该显示一条消息警告用户。如果进行了选择,它将导航到另
我知道 MSVC 可以通过 pragma 消息做到这一点 -> http://support.microsoft.com/kb/155196 gcc 是否有办法打印用户创建的警告或消息? (我找不到谷
当存在大量节点或二进制数据时, native Erlang 消息能否提供合理的性能? 情况 1:有一个大约 50-200 台机器的动态池(erlang 节点)。它在不断变化,每 10 分钟大约添加或删
我想知道如何在用户登录后显示“欢迎用户,您已登录”的问候消息,并且该消息应在 5 秒内消失。 该消息将在用户成功登录后显示一次,但在同一 session 期间连续访问主页时不会再次显示。因为我在 ho
如果我仅使用Welcome消息,我的代码可以正常工作,但是当打印p->client_name指针时,消息不居中。 所以我的问题是如何将消息和客户端名称居中,就像它是一条消息一样。为什么它目前仅将消
我是一名优秀的程序员,十分优秀!