如何从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
相关产品推荐
相关产品推荐

