Skip to content

Latest commit

History

41 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

python-stream

说明

数据流式框架, 可用作数据清洗, 数据预处理, 数据迁移等应用场景

更优雅的流式数据处理方式

安装


pip install git+https://github.com/sandabuliu/python-stream.git

or

git clone https://github.com/sandabuliu/python-stream.git
cd python-agent
python setup.py install

QuickStart


Examples

Word Count
frompystream.executor.sourceimportMemoryfrompystream.executor.executorimportMap, Iterator, ReducebyKeydata=Memory([
'Wikipedia is a free online encyclopedia, created and edited by volunteers around the world and hosted by the Wikimedia Foundation.',
'Search thousands of wikis, start a free wiki, compare wiki software.',
'The official Wikipedia Android app is designed to help you find, discover, and explore knowledge on Wikipedia.'
])
p=data|Map(lambdax: x.split(' ')) |Iterator(lambdax: (x.strip('.,'), 1)) |ReducebyKey(lambdax, y: x+y)
result= {}
forkey, valueinp:
result[key] =valueprintresult.items()

执行结果

[('and', 3), ('wiki', 2), ('compare', 1), ('help', 1), ('is', 2), ('Wikipedia', 3), ('discover', 1), ('hosted', 1), ('Android', 1), ('find', 1), ('Foundation', 1), ('knowledge', 1), ('to', 1), ('by', 2), ('start', 1), ('online', 1), ('you', 1), ('thousands', 1), ('app', 1), ('edited', 1), ('Search', 1), ('around', 1), ('free', 2), ('explore', 1), ('designed', 1), ('world', 1), ('The', 1), ('the', 2), ('a', 2), ('on', 1), ('created', 1), ('Wikimedia', 1), ('official', 1), ('encyclopedia', 1), ('of', 1), ('wikis', 1), ('volunteers', 1), ('software', 1)]
计算π
fromrandomimportrandomfrompystream.executor.sourceimportFakerfrompystream.executor.executorimportExecutor, Map, GroupclassPi(Executor):
def__init__(self, **kwargs):
super(Pi, self).__init__(**kwargs)
self.counter=0self.result=0defhandle(self, item):
self.counter+=1self.result+=itemreturn4.0*self.result/self.counters=Faker(lambda: random(), 100000) |Map(lambdax: x*2-1) |Group(size=2) |Map(lambdax: 1ifx[0]**2+x[1]**2<=1else0) |Pi()
res=Nonefor_ins:
res=_printres

执行结果

3.14728
排序
fromrandomimportrandintfrompystream.executor.sourceimportMemoryfrompystream.executor.executorimportSortm=Memory([randint(0, 100) foriinrange(10)]) |Sort()
foriinm:
printlist(i)

执行结果

[94]
[94, 99]
[18, 94, 99]
[18, 40, 94, 99]
[18, 26, 40, 94, 99]
[18, 26, 40, 63, 94, 99]
[18, 26, 40, 63, 83, 94, 99]
[3, 18, 26, 40, 63, 83, 94, 99]
[3, 18, 26, 40, 63, 83, 83, 94, 99]
[3, 16, 18, 26, 40, 63, 83, 83, 94, 99]
在 hadoop 中使用
wordcount
mapper.py
frompystream.executor.sourceimportStdinfrompystream.executor.executorimportMap, Iteratorfrompystream.executor.outputimportStdouts=Stdin() |Map(lambdax: x.strip().split()) |Iterator(lambdax: "%s\t1"%x) |Stdout()
s.start()
reducer.py
frompystream.executor.sourceimportStdinfrompystream.executor.executorimportMap, ReducebySortedKeyfrompystream.executor.outputimportStdouts=Stdin() |Map(lambdax: x.strip().split('\t')) |ReducebySortedKey(lambdax, y: int(x)+int(y)) |Map(lambdax: '%s\t%s'%x) |Stdout()
s.start()
解析 NGINX 日志
frompystream.configimportrulefrompystream.executor.sourceimportFilefrompystream.executor.executorimportParsers=File('/var/log/nginx/access.log') |Parser(rule('nginx'))
foritemins:
printitem

执行结果

{'status': '400', 'body_bytes_sent': 173, 'remote_user': '-', 'http_referer': '-', 'remote_addr': '198.35.46.20', 'request': '\\x05\\x01\\x00', 'version': None, 'http_user_agent': '-', 'time_local': datetime.datetime(2017, 2, 15, 13, 11, 3), 'path': None, 'method': None}
{'status': '400', 'body_bytes_sent': 173, 'remote_user': '-', 'http_referer': '-', 'remote_addr': '198.35.46.20', 'request': '\\x05\\x01\\x00', 'version': None, 'http_user_agent': '-', 'time_local': datetime.datetime(2017, 2, 15, 13, 11, 3), 'path': None, 'method': None}
{'status': '400', 'body_bytes_sent': 173, 'remote_user': '-', 'http_referer': '-', 'remote_addr': '198.35.46.20', 'request': '\\x05\\x01\\x00', 'version': None, 'http_user_agent': '-', 'time_local': datetime.datetime(2017, 2, 15, 13, 11, 3), 'path': None, 'method': None}
{'status': '400', 'body_bytes_sent': 173, 'remote_user': '-', 'http_referer': '-', 'remote_addr': '198.35.46.20', 'request': '\\x05\\x01\\x00', 'version': None, 'http_user_agent': '-', 'time_local': datetime.datetime(2017, 2, 15, 13, 11, 3), 'path': None, 'method': None}
{'status': '400', 'body_bytes_sent': 173, 'remote_user': '-', 'http_referer': '-', 'remote_addr': '198.35.46.20', 'request': '\\x05\\x01\\x00', 'version': None, 'http_user_agent': '-', 'time_local': datetime.datetime(2017, 2, 15, 13, 11, 3), 'path': None, 'method': None}
导出数据库数据
fromsqlalchemyimportcreate_enginefrompystream.executor.sourceimportSQLfrompystream.executor.outputimportCsvfrompystream.executor.wrapsimportBatchengine=create_engine('mysql://root:123456@127.0.0.1:3306/test') conn=engine.connect()
s=SQL(conn, 'select * from faker') |Batch(Csv('/tmp/output'))
foritemins:
printitem['data']
printitem['exception']
conn.close()

数据源

读取文件数据
frompystream.executor.sourceimportTail, File, CsvTail('/var/log/nginx/access.log')
File('/var/log/nginx/*.log')
Csv('/tmp/test*.csv')
读取 TCP 流数据
frompystream.executor.sourceimportTCPClientTCPClient('/tmp/pystream.sock')
TCPClient(('127.0.0.1', 10000))
读取 python 数据
fromQueueimportQueueasQfromrandomimportrandintfrompystream.executor.sourceimportMemory, Faker, Queuequeue=Q(10)
Memory([1, 2, 3, 4])
Faker(randint, 1000)
Queue(queue)
读取常用模块数据
frompystream.executor.sourceimportSQL, KafkaSQL(conn, 'select * from faker') # 读取数据库数据Kafka('topic1', '127.0.0.1:9092') # 读取 kafka 数据

数据输出

输出到文件
frompystream.executor.outputimportFile, CsvFile('/tmp/output')
Csv('/tmp/output.csv')
通过HTTP输出
frompystream.executor.outputimportHTTPRequestHTTPRequest('http://127.0.0.1/api/data')
输出到kafka
frompystream.executor.outputimportKafkaKafka('topic', '127.0.0.1:9092')

中间件

队列
frompystream.executor.sourceimportTailfrompystream.executor.outputimportStdoutfrompystream.executor.middlewareimportQueues=Tail('/Users/tongbin01/PycharmProjects/python-stream/README.md') |Queue() |Stdout()
s.start()
订阅
fromrandomimportrandintfrompystream.executor.sourceimportTailfrompystream.executor.executorimportMapfrompystream.executor.outputimportStdoutfrompystream.executor.middlewareimportSubscribefrompystream.executor.wrapsimportDaemonicsub=Tail('/var/log/messages') |Map(lambdax: (str(randint(1, 2)), x.strip())) |Subscribe()
Daemonic(sub).start()
s=sub['1'] |Map(lambdax: x.strip()) |Stdout()
s.start()

TodoList

  • 订阅器(Subscribe)客户端超时处理
  • 并行计算
  • HTTP 异步输出/异步源
  • 添加其他基础输出/基础源
  • 添加对其他常用模块的支持, 如 redis, kafka, flume, log-stash, 各种数据库等

Copyright © 2017 g_tongbin@foxmail.com

About

更优雅的流式数据处理方式

Topics

Resources

Stars

31 stars

Watchers

3 watching

Forks

Releases

Packages

Used by

Contributors

Languages