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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 09:05:10