- 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的文章或继续浏览相关文章,希望大家以后支持我的博客! 。
我正在处理一组标记为 160 个组的 173k 点。我想通过合并最接近的(到 9 或 10 个组)来减少组/集群的数量。我搜索过 sklearn 或类似的库,但没有成功。 我猜它只是通过 knn 聚类
我有一个扁平数字列表,这些数字逻辑上以 3 为一组,其中每个三元组是 (number, __ignored, flag[0 or 1]),例如: [7,56,1, 8,0,0, 2,0,0, 6,1,
我正在使用 pipenv 来管理我的包。我想编写一个 python 脚本来调用另一个使用不同虚拟环境(VE)的 python 脚本。 如何运行使用 VE1 的 python 脚本 1 并调用另一个 p
假设我有一个文件 script.py 位于 path = "foo/bar/script.py"。我正在寻找一种在 Python 中通过函数 execute_script() 从我的主要 Python
这听起来像是谜语或笑话,但实际上我还没有找到这个问题的答案。 问题到底是什么? 我想运行 2 个脚本。在第一个脚本中,我调用另一个脚本,但我希望它们继续并行,而不是在两个单独的线程中。主要是我不希望第
我有一个带有 python 2.5.5 的软件。我想发送一个命令,该命令将在 python 2.7.5 中启动一个脚本,然后继续执行该脚本。 我试过用 #!python2.7.5 和http://re
我在 python 命令行(使用 python 2.7)中,并尝试运行 Python 脚本。我的操作系统是 Windows 7。我已将我的目录设置为包含我所有脚本的文件夹,使用: os.chdir("
剧透:部分解决(见最后)。 以下是使用 Python 嵌入的代码示例: #include int main(int argc, char** argv) { Py_SetPythonHome
假设我有以下列表,对应于及时的股票价格: prices = [1, 3, 7, 10, 9, 8, 5, 3, 6, 8, 12, 9, 6, 10, 13, 8, 4, 11] 我想确定以下总体上最
所以我试图在选择某个单选按钮时更改此框架的背景。 我的框架位于一个类中,并且单选按钮的功能位于该类之外。 (这样我就可以在所有其他框架上调用它们。) 问题是每当我选择单选按钮时都会出现以下错误: co
我正在尝试将字符串与 python 中的正则表达式进行比较,如下所示, #!/usr/bin/env python3 import re str1 = "Expecting property name
考虑以下原型(prototype) Boost.Python 模块,该模块从单独的 C++ 头文件中引入类“D”。 /* file: a/b.cpp */ BOOST_PYTHON_MODULE(c)
如何编写一个程序来“识别函数调用的行号?” python 检查模块提供了定位行号的选项,但是, def di(): return inspect.currentframe().f_back.f_l
我已经使用 macports 安装了 Python 2.7,并且由于我的 $PATH 变量,这就是我输入 $ python 时得到的变量。然而,virtualenv 默认使用 Python 2.6,除
我只想问如何加快 python 上的 re.search 速度。 我有一个很长的字符串行,长度为 176861(即带有一些符号的字母数字字符),我使用此函数测试了该行以进行研究: def getExe
list1= [u'%app%%General%%Council%', u'%people%', u'%people%%Regional%%Council%%Mandate%', u'%ppp%%Ge
这个问题在这里已经有了答案: Is it Pythonic to use list comprehensions for just side effects? (7 个答案) 关闭 4 个月前。 告
我想用 Python 将两个列表组合成一个列表,方法如下: a = [1,1,1,2,2,2,3,3,3,3] b= ["Sun", "is", "bright", "June","and" ,"Ju
我正在运行带有最新 Boost 发行版 (1.55.0) 的 Mac OS X 10.8.4 (Darwin 12.4.0)。我正在按照说明 here构建包含在我的发行版中的教程 Boost-Pyth
学习 Python,我正在尝试制作一个没有任何第 3 方库的网络抓取工具,这样过程对我来说并没有简化,而且我知道我在做什么。我浏览了一些在线资源,但所有这些都让我对某些事情感到困惑。 html 看起来
我是一名优秀的程序员,十分优秀!