使用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
相关产品推荐
相关产品推荐

