如何使用Dask直接读取HDFS中的sas7bdat文件?
解决dask_sas_reader读取HDFS文件时路径解析错误问题
在YARN集群上使用Dask + dask_sas_reader读取HDFS中的sas7bdat文件时,出现路径解析错误:工具将hdfs://dev/data/airline.sas7bdat错误拼接为/root/hdfs:/dev/data/airline.sas7bdat,导致找不到文件,同时触发IndexError。
问题原因
dask_sas_reader依赖的pyreadstat是面向本地文件系统的工具,不直接支持HDFS路径。当传入HDFS格式的路径时,工具会默认将其视为本地路径的一部分,与当前工作目录(/root)拼接,从而生成无效路径。
解决方案
由于dask_sas_reader原生不支持HDFS,需要通过Dask的HDFS文件系统接口先处理文件,再传递给pyreadstat处理。以下提供两种可行方案:
方案一:将HDFS文件同步到Worker本地临时目录
先将HDFS上的sas7bdat文件复制到每个YARN Worker的本地临时目录,再用dask_sas_reader读取本地文件:
from dask_sas_reader import sas from dask.distributed import Client from dask_yarn import YarnCluster from dask.hdfs import HDFS import tempfile import os # 部署YARN集群并连接 cluster = YarnCluster(environment='demo.tar.gz') cluster.scale(2) client = Client(cluster) # 初始化HDFS客户端,定义路径 hdfs = HDFS() hdfs_file_path = "/dev/data/airline.sas7bdat" local_temp_path = os.path.join(tempfile.gettempdir(), "airline.sas7bdat") # 定义Worker端执行的文件复制函数 def copy_hdfs_to_local(hdfs_path, local_path): with hdfs.open(hdfs_path, 'rb') as hdfs_file: with open(local_path, 'wb') as local_file: local_file.write(hdfs_file.read()) return local_path # 在所有Worker上同步文件 client.run(copy_hdfs_to_local, hdfs_file_path, local_temp_path) # 读取本地文件并转换 dd_df = sas.dask_sas_reader(local_temp_path, blocksize=800000) dd_df.compute().to_parquet("/root/airline.parquet") # 清理Worker本地临时文件 client.run(os.remove, local_temp_path) # 关闭集群连接 client.shutdown() cluster.shutdown()
方案二:基于Dask字节流分块处理
直接读取HDFS文件的字节块,封装自定义读取逻辑适配pyreadstat:
from dask.distributed import Client from dask_yarn import YarnCluster from dask.bytes import read_bytes from dask.dataframe import from_delayed import pyreadstat import tempfile import os # 部署YARN集群并连接 cluster = YarnCluster(environment='demo.tar.gz') cluster.scale(2) client = Client(cluster) # 读取HDFS文件的分块字节流 blocks = read_bytes("hdfs://dev/data/airline.sas7bdat", blocksize=800000) # 定义单块读取函数:字节流转临时文件后用pyreadstat读取 def read_sas_block(block_data): with tempfile.NamedTemporaryFile(suffix='.sas7bdat', delete=False) as tmp_file: tmp_file.write(block_data) tmp_path = tmp_file.name df, _ = pyreadstat.read_sas7bdat(tmp_path) os.unlink(tmp_path) # 清理临时文件 return df # 生成延迟任务并构建Dask DataFrame delayed_dfs = [from_delayed(read_sas_block(block)) for block in blocks] dd_df = from_delayed(delayed_dfs) # 计算并保存结果 dd_df.compute().to_parquet("/root/airline.parquet") # 关闭集群连接 client.shutdown() cluster.shutdown()
注意事项
- 确保YARN Worker节点的临时目录有足够空间存放sas7bdat文件;
- 方案二适合大文件分块处理,避免单个Worker加载整个文件;
- 环境包
demo.tar.gz需包含pyreadstat、dask-sas-reader及HDFS相关依赖(如hdfs3)。
内容的提问来源于stack exchange,提问作者aerylias
相关产品推荐
相关产品推荐

