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

如何在Kedro中将Azure Databricks Lakehouse查询纳入DataCatalog数据集?

Kedro管控Azure Databricks机器学习流水线的解决方案

核心问题解答

  1. 能否将Azure Databricks Lakehouse的查询作为Kedro DataCatalog数据集?
    可以。Kedro支持通过多种数据集类型对接Databricks Lakehouse,根据数据规模和技术栈选择即可:
  • 若使用Spark(Databricks原生引擎),优先用pyspark.SQLQueryDataSet,完全适配分布式大数据场景;
  • 若需用Pandas处理,可使用pandas.SQLQueryDataSet,搭配Databricks SQL连接器实现查询。
  1. Databricks中能否实现内存高效的查询关联?
    完全可以。Databricks基于Spark构建,天然具备以下特性保障内存高效:
  • 懒执行机制:Spark SQL查询仅在触发动作(如写入结果、展示数据)时才执行,不会提前加载全表到内存;
  • 智能优化:Spark Catalyst优化器会自动做谓词下推、分区裁剪、关联顺序优化,只处理符合查询条件的数据集;
  • 分布式处理:数据分散在集群节点处理,单节点无需加载全表,避免内存溢出。

配置示例

1. 适配大数据场景的PySpark配置(推荐)

使用pyspark.SQLQueryDataSet对接Databricks Lakehouse,天然支持内存高效的关联查询:

scooters_joined_data:
  type: pyspark.SQLQueryDataSet
  sql: |
    SELECT c.*, s.scooter_id, s.location
    FROM cars c
    INNER JOIN scooters s ON c.owner_id = s.owner_id
    WHERE c.gear = 4
  load_args:
    cache: false  # 无需重复使用结果时关闭缓存,减少内存占用
  credentials: databricks_spark_creds
  • databricks_spark_creds需配置SparkSession连接信息(可通过Kedro hooks初始化Databricks SparkSession);
  • 该配置下,Spark会自动利用Delta Lake表的分区信息裁剪数据,仅扫描gear=4的分区,关联操作也会分布式执行,不会加载全表到内存。

2. Pandas场景的内存优化配置

若需用Pandas处理,可通过分块加载避免一次性加载全表:

scooters_query:
  type: pandas.SQLQueryDataSet
  credentials: databricks_sql_creds
  sql: SELECT * FROM cars WHERE gear = 4
  load_args:
    index_col: [name]
    chunksize: 10000  # 按1万行分块加载,降低单批次内存占用
  • databricks_sql_creds需包含Databricks SQL的连接信息(server_hostname、http_path、access_token等);
  • chunksize参数控制每次加载的数据量,适合处理大表但结果集仍可单节点处理的场景。

关键优化建议

  • 优先选择PySpark数据集类型,适配Databricks的分布式架构,最大化内存效率;
  • 编写SQL时尽量过滤掉不必要的数据,利用Lakehouse表的分区、索引特性,进一步减少扫描的数据量;
  • 避免对大结果集做全量缓存,仅在重复使用数据时开启Spark缓存。

内容的提问来源于stack exchange,提问作者Jon Nagra

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 17:20:33