Issues scaling kafka consumers with Dask
还没有人认领这个 Issue。
评估
- 难度
- 5/5
- 预计耗时
- 一周以上
- 新手友好度
- 25/100
- Issue 类型
- 缺陷
- 描述清晰度
- 需要澄清
- 活跃度
- 停滞
- 技术栈
- kafka, python
调研方向
Start by running the supplied Stream.from_kafka_batched pipelines with dask disabled and enabled, varying Kafka topic partitions as described. Inspect the LocalCluster configuration and compare throughput measurements; done means documenting the cause of the scaling behavior and identifying a concrete corrective change, since no file or test is named.
由索引模型根据 Issue 内容生成。
描述
I have run two simple pipelines(Non-dask and Dask) to read stream of messages from Kafka using Streamz FromKafkaBatched method. I seem to get similar throughput for both the pipelines.
Non-Dask pipeline:
from distributed import Client
from streamz import Stream
import time
import json
def preprocess(messages):
start_time = int(round(time.time()))
no_of_rows = len(messages)
if no_of_rows > 0:
size = no_of_rows*len(messages[0])
kafka_timestamp = json.loads(messages[0].decode('utf-8'))['timestamp']
else:
size = 0
kafka_timestamp = start_time
end_time = int(round(time.time()))
return "{},{},{},{},{}".format(no_of_rows, kafka_timestamp, start_time, end_time, size)
topic = "topic-1"
bootstrap_servers = 'localhost:9092'
consumer_conf = {'bootstrap.servers': bootstrap_servers,
'group.id': 'group_1', 'session.timeout.ms': 60000}
stream = Stream.from_kafka_batched(topic, consumer_conf, poll_interval='10s',
npartitions=5, asynchronous = True, dask= False)
kafka_out = stream.map(preprocess).to_kafka('test-out', consumer_conf)
kafka_out.flush()
The pipeline gives maximum throughput = ~43 MBps (topic partitions = 5).
Throughput computation = (Σ size) / (max(end_time) - min(start_time))
By setting dask=True the above pipeline is executed on a local dask cluster
Starting Dask Cluster:
# dask imports
from distributed import Client, LocalCluster
# dask client
cluster = LocalCluster(ip="0.0.0.0", scheduler_port=8786, diagnostics_port = 8787, processes=False, threads_per_worker=10, n_workers=1)
client = Client(cluster)
client.get_versions(check=True)
The pipeline with Dask gives maximum throughput of ~44.8 MBps when number of partitions in the topic = 1. Increasing the number of partitions in a topic decreases the performance.
My understanding from these experimentations is that confluent kafka performs as expected with single Dask thread (when number of partitions is 1) but performs poorly with multiple dask threads or multiple partitions. This may be due to underlying tornado async library. It would be great if experts help me understand this better as I am fairly new to async programming.
Hardware and Software Specs:
10 CPU cores, 30 GB RAM
Ubuntu 16.04, Streamz(latest-build off latest code base)
Kafka 2.2.0, confluent-kafka-python 1.0.0
- 主要语言
- Python
- 星标
- 1.3k
- 派生
- 150
- 平均合并
- 17 小时 39 分钟
- 30 天内合并 PR
- 1
环境准备
- 提供 Dockerfile 或 Docker Compose 文件
- 没有 Pull Request 模板
- 阅读贡献指南
从这里开始
- 先读完整个 Issue,再读项目的贡献指南。
- 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
- Fork 仓库,在一个分支上完成修改。
- 提交 Pull Request,并在描述里引用这个 Issue 编号。
python-streamz/streamz 的其他 Issue
-
难度 3/5 1-2 天 新手友好度 35/100
python-streamz/streamz#481 · 4 条评论 ·
-
难度 5/5 一周以上 新手友好度 25/100
python-streamz/streamz#480 · 2 条评论 ·
-
难度 5/5 一周以上 新手友好度 20/100
python-streamz/streamz#479 · 1 条评论 ·
-
难度 5/5 一周以上 新手友好度 20/100
python-streamz/streamz#478 · 6 条评论 ·
-
难度 4/5 3-5 天 新手友好度 10/100
python-streamz/streamz#476 · 17 条评论 · 2 个 reaction ·
查看 python-streamz/streamz 的全部 Issue
相似的 Issue
-
Maven path-index: "Ambiguous or noncanonical artifact path" error does not report the offending path未关闭
难度 2/5 1-3 小时 新手友好度 76/100
pulp/pulp_maven#524 ·
维护者通常 1 天内回复
-
难度 1/5 1-3 小时 新手友好度 82/100
维护者通常 1 天内回复
-
难度 2/5 1-3 小时 新手友好度 76/100
维护者通常 3 天内回复
-
难度 1/5 1-3 小时 新手友好度 88/100
infinispan/langchain-infinispan#34 ·
维护者通常 1 天内回复
-
难度 1/5 1 小时以内 新手友好度 85/100
521xueweihan/HelloGitHub#3891 ·