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

如何从Data Lake Blob Storage导入非结构化CSV至Databricks并读取全量内容

问题描述

我正尝试从Data Lake存储将非结构化CSV文件导入Databricks,并希望读取该文件的全部内容,文件内容如下:

EdgeMaster                          
Name    Value   Unit    Status  Nom.    Lower   Upper   Description
Type    A                                                  A
Date    1/1/2022                                           B
Time    0:00:00                                            A
X       1       m       OK      1       2       3          B
Y       -                                                  A
                            
EdgeMaster                          
Name    Value   Unit    Status  Nom.    Lower   Upper   Description
Type    B                                                  C
Date    1/1/2022                                           D
Time    0:00:00                                            C
X       1       m       OK      1       2       3          D
Y       -                                                  C

我尝试了以下三种方法,但均未成功读取全量内容:

  • 方法1:将首行作为表头读取
    df = sqlContext.read.format("com.databricks.spark.csv").option("header", "true").load('abfss://xyz/sample.csv')
    
    仅得到部分结果。
  • 方法2:跳过表头读取,结果无改善。
  • 方法3:定义自定义Schema,查询未返回任何结果。
解决方案

这个CSV的特殊之处在于它包含重复的数据集块,每个块以EdgeMaster开头且自带表头,普通CSV读取器无法识别这种嵌套结构,需要手动按行解析:

1. 读取所有行并预处理

先把文件当作纯文本读取,过滤空行并保留原始内容:

# 读取所有行,过滤空行
raw_lines = spark.read.text("abfss://xyz/sample.csv").filter("value != ''").rdd.map(lambda x: x[0])

2. 解析数据块

遍历所有行,识别每个EdgeMaster开头的数据块,将块内的数据映射到对应表头:

from pyspark.sql import Row

result_rows = []
current_block_data = []
current_header = None

for line in raw_lines.collect():
    # 按制表符分割行内容(原文件是制表分隔),清理空元素
    split_line = [item.strip() for item in line.split("\t") if item.strip()]
    
    if line.startswith("EdgeMaster"):
        # 遇到新块,先处理上一个块的内容
        if current_header and current_block_data:
            for data_line in current_block_data:
                # 把数据行和表头映射成字典,再转成Row对象
                row_dict = dict(zip(current_header, data_line))
                result_rows.append(Row(**row_dict))
        # 重置当前块的状态
        current_header = None
        current_block_data = []
    elif current_header is None:
        # 该行是当前块的表头
        current_header = split_line
    else:
        # 该行是当前块的数据行,加入临时列表
        current_block_data.append(split_line)

# 处理最后一个未完成的数据块
if current_header and current_block_data:
    for data_line in current_block_data:
        row_dict = dict(zip(current_header, data_line))
        result_rows.append(Row(**row_dict))

# 转换为DataFrame
final_df = spark.createDataFrame(result_rows)
final_df.show()

关键说明

  • 原文件使用制表符分隔字段,而非逗号,所以用split("\t")分割行内容
  • 过滤空行和分割后的空元素,避免字段错位
  • 每个EdgeMaster标记一个独立数据集,我们逐个解析并合并所有数据行
  • 最终的DataFrame会包含所有块的完整数据,不会丢失任何行

内容的提问来源于stack exchange,提问作者Ankit Sawa

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 18:31:49