如何在Kedro中将Azure Databricks Lakehouse查询纳入DataCatalog数据集?
Kedro管控Azure Databricks机器学习流水线的解决方案
核心问题解答
- 能否将Azure Databricks Lakehouse的查询作为Kedro DataCatalog数据集?
可以。Kedro支持通过多种数据集类型对接Databricks Lakehouse,根据数据规模和技术栈选择即可:
- 若使用Spark(Databricks原生引擎),优先用
pyspark.SQLQueryDataSet,完全适配分布式大数据场景; - 若需用Pandas处理,可使用
pandas.SQLQueryDataSet,搭配Databricks SQL连接器实现查询。
- 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
相关产品推荐
相关产品推荐

