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

使用Apache Beam read_csv读取HTTP在线CSV文件失败求助

解决Apache Beam read_csv读取HTTP URL失败的问题

问题根源

apache_beam.dataframe.io.read_csv依赖Beam的FileSystem抽象,目前官方只支持本地文件、GS、S3等文件系统协议,不直接支持HTTP/HTTPS网络URL,所以会抛出无法获取文件系统的错误。

解决思路

方法1:先下载HTTP文件到本地/临时存储,再读取

手动将在线CSV下载到本地临时文件,再用read_csv读取本地路径,这是最直接的方案:

import apache_beam as beam
from apache_beam.dataframe.io import read_csv
import requests
import tempfile
import os

url = 'https://github.com/datablist/sample-csv-files/raw/main/files/people/people-100.csv'

# 下载文件到临时文件
response = requests.get(url)
with tempfile.NamedTemporaryFile(mode='w+b', suffix='.csv', delete=False) as temp_file:
    temp_file.write(response.content)
    temp_path = temp_file.name

try:
    with beam.Pipeline() as pipeline:
        df = pipeline | read_csv(path=temp_path)
        df[:5] | beam.Map(lambda row: print(row))
finally:
    # 清理临时文件
    os.unlink(temp_path)

方法2:通过Beam原生PCollection读取并解析HTTP CSV

绕过DataFrame API的read_csv,直接用Beam的原生操作读取HTTP内容、解析CSV,再转换成DataFrame格式:

import apache_beam as beam
import csv
from io import StringIO

url = 'https://github.com/datablist/sample-csv-files/raw/main/files/people/people-100.csv'

def download_and_parse_csv(url):
    import requests
    response = requests.get(url)
    content = response.text
    reader = csv.DictReader(StringIO(content))
    return list(reader)

with beam.Pipeline() as pipeline:
    # 下载并解析CSV为字典列表
    csv_rows = pipeline | beam.Create([url]) | beam.FlatMap(download_and_parse_csv)
    # 转换成Beam DataFrame
    df = csv_rows | beam.dataframe.convert.to_dataframe()
    # 处理并输出前5行
    df[:5] | beam.Map(lambda row: print(row))

补充说明

如果需要频繁处理HTTP数据源,可以考虑封装自定义的FileSystem实现,但官方目前没有提供HTTP的FileSystem插件,所以前两种方法是更实用的选择。

内容的提问来源于stack exchange,提问作者Nikolai Polikurov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 16:15:54