使用dbt以Spark为引擎构建湖仓:如何读取JSON生成Delta/Iceberg表?
问题解答
dbt-spark确实没有内置的原生函数(类似dbt-duckdb的read_json/read_csv)直接读取原始文件并写入Delta/Iceberg表,但完全可以在dbt框架内通过以下几种方式实现,无需单独编写Spark作业:
方法1:直接使用Spark SQL创建表并导入数据
在dbt模型中编写Spark原生SQL,通过CREATE TABLE ... AS SELECT语法读取原始文件并生成Delta/Iceberg表,dbt仅作为SQL执行载体。
示例:生成Delta表
-- models/raw/my_raw_delta_table.sql CREATE OR REPLACE TABLE my_raw_delta_table USING delta LOCATION 's3://your-bucket/delta-tables/raw-data' -- 指定Delta表存储路径 AS SELECT * FROM json.`s3://your-bucket/raw/json-files/` -- 读取原始JSON文件
示例:生成Iceberg表
-- models/raw/my_raw_iceberg_table.sql CREATE OR REPLACE TABLE my_raw_iceberg_table USING iceberg LOCATION 's3://your-bucket/iceberg-tables/raw-data' -- 指定Iceberg表存储路径 AS SELECT * FROM json.`s3://your-bucket/raw/json-files/`
将这类SQL文件放在dbt的models/raw目录下,执行dbt run即可完成数据导入和表创建。
方法2:自定义dbt宏实现灵活导入
如果需要动态路径、分区处理等复杂逻辑,可以编写自定义宏,通过dbt的spark_session对象调用Spark DataFrame API完成导入。
示例宏定义(放在macros/load_raw_data.sql)
{% macro load_json_to_format(source_path, target_table, target_location, table_format='delta') %} {% set spark = spark_session %} -- 读取JSON文件 {% set df = spark.read.option("multiline", "true").json(source_path) %} -- 写入目标格式表 {% do df.write .format(table_format) .mode("overwrite") .option("path", target_location) .saveAsTable(target_table) %} {% endmacro %}
执行宏命令
dbt run-operation load_json_to_format --args '{ "source_path": "s3://your-bucket/raw/json-files/", "target_table": "my_raw_table", "target_location": "s3://your-bucket/tables/raw-data", "table_format": "iceberg" }'
方法3:通过dbt外部表配置指向原始文件(仅映射,不复制数据)
如果仅需要创建指向原始JSON文件的外部表,后续再转换为Delta/Iceberg,可以使用dbt-spark的external物化配置:
-- models/raw/my_json_external_table.sql {{ config( materialized='external', location='s3://your-bucket/raw/json-files/', format='json' ) }} SELECT * FROM {{ source('raw_sources', 'json_data') }}
其中raw_sources需在sources.yml中定义:
sources: - name: raw_sources tables: - name: json_data external: true location: 's3://your-bucket/raw/json-files/' format: 'json'
后续若要将外部表转换为Delta/Iceberg,只需基于该外部表执行CREATE TABLE ... AS SELECT即可。
内容的提问来源于stack exchange,提问作者nerdizzle
相关产品推荐
相关产品推荐

