如何在Databricks中用PySpark从Azure容器动态读取配置指定列的CSV
实现步骤
1. 配置config.ini文件
先创建config.ini文件,指定需要加载的目标列:
[CSV_SETTINGS] target_column = Empid
2. 在Databricks中读取配置
用Python的configparser库读取配置文件,注意路径要对应文件实际存储位置(比如Databricks工作区路径或挂载的Azure存储路径):
import configparser # 加载配置文件 config = configparser.ConfigParser() config.read("/dbfs/FileStore/config.ini") # 替换为你的config.ini实际路径 # 提取指定列 target_col = config.get("CSV_SETTINGS", "target_column")
3. 读取Azure容器中的CSV并处理列
先确保Azure容器已挂载到Databricks(比如通过ADLS Gen2挂载),读取CSV时仅加载目标列,剩余列填充null值:
from pyspark.sql.functions import lit # CSV文件的挂载路径 csv_file_path = "/dbfs/mnt/adls_container/employee.csv" # Delta目标表的列名列表 delta_table_cols = ["Empid", "Ename", "Esalary"] # 读取CSV,仅加载指定列 raw_df = spark.read.options( header=True, inferSchema=True, columns=[target_col] ).csv(csv_file_path) # 复用Delta表的schema生成null列,保证类型匹配 delta_schema = spark.read.table("your_delta_target_table").schema for field in delta_schema.fields: col_name = field.name if col_name not in raw_df.columns: raw_df = raw_df.withColumn(col_name, lit(None).cast(field.dataType)) # 调整列顺序与Delta表一致 final_df = raw_df.select(delta_table_cols)
4. 写入Delta Lake目标表
根据业务需求选择写入模式(覆盖/追加):
# 写入Delta表 final_df.write.mode("overwrite").saveAsTable("your_delta_target_table")
额外说明
- 如果需要支持多列配置,只需把
config.ini中的target_column改成逗号分隔格式(比如target_column = Empid,Ename),再用target_cols = config.get("CSV_SETTINGS", "target_column").split(",")提取列列表即可 - 可通过
dbutils.fs.ls("/dbfs/mnt/adls_container")验证Azure容器挂载路径的正确性 - 复用Delta表schema生成null列能避免写入时的类型不兼容问题
内容的提问来源于stack exchange,提问作者bigdata techie
相关产品推荐
相关产品推荐

