You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.05 02:06:03