如何让Dask DataFrame从标准输入stdin读取数据
Dask DataFrame 从标准输入读取数据的解决方案
方法1:类Unix平台直接使用/dev/stdin路径
Linux、macOS等类Unix系统中存在/dev/stdin虚拟路径,对应当前进程的标准输入流,可直接传入Dask的read_csv方法使用:
import dask.dataframe as dd # 直接传入标准输入虚拟路径,指定blocksize=None禁用分块读取 df = dd.read_csv("/dev/stdin", sep=" ", header=None, blocksize=None) # 计算时指定同步调度器,避免多进程重复读取标准输入 result = df.compute(scheduler="synchronous")
- 该方案代码改动最小,缺点是不支持Windows平台
- 必须指定
blocksize=None,因为标准输入是不可寻址的流,不支持随机分块读取
方法2:跨平台通用方案:通过dask.delayed手动封装读取逻辑
全平台通用,无需将数据落地到本地文件,自己控制分块读取标准输入后拼接为Dask DataFrame:
import dask import dask.dataframe as dd import pandas as pd import sys from io import StringIO def parse_chunk(lines): buffer = StringIO("\n".join(lines)) return pd.read_csv(buffer, sep=" ", header=None) # 可根据内存情况调整单块行数,平衡调度开销和内存占用 chunk_size = 10000 chunk_list = [] current_chunk = [] for line in sys.stdin: current_chunk.append(line.strip()) if len(current_chunk) >= chunk_size: chunk_list.append(dask.delayed(parse_chunk)(current_chunk)) current_chunk = [] # 处理最后不足块大小的剩余数据 if current_chunk: chunk_list.append(dask.delayed(parse_chunk)(current_chunk)) # 拼接所有分块为完整的Dask DataFrame df = dd.from_delayed(chunk_list)
- 读取标准输入的阶段为串行执行,符合标准输入仅能顺序读取一次的特性,后续的计算环节可正常并行
关键注意事项
- 标准输入是不可寻址、仅能顺序读取一次的数据流,任何方案都无法实现多进程并行读取标准输入
- 不要在分布式Dask集群中使用标准输入作为数据源,集群worker无法访问主进程的标准输入流
内容的提问来源于stack exchange,提问作者Win
相关产品推荐
相关产品推荐

