Snowflake-Pyspark配置numPartitions未生效问题咨询
问题原因及解决方案
1. 核心触发原因
你当前遇到的单条查询问题,是因为Spark原生JDBC的分区拆分逻辑仅在使用dbtable参数指定表/子查询时生效,你代码中使用了query参数传入自定义SQL,Spark不会自动做查询拆分,只会提交1次全量查询,因此Snowflake后台只会出现1条查询记录。
此外如果你使用的是Snowflake官方提供的Spark连接器,配置规则和原生JDBC存在差异,也会导致分区参数不生效。
2. 修复方案
方案一:保留原生JDBC用法,调整配置项
将自定义query参数替换为dbtable参数即可触发分区逻辑,修改后代码示例:
sfOptions = dict() sfOptions["url"] ="jdbc:snowflake://**************.privatelink.snowflakecomputing.com" sfOptions["user"] ="**01d" sfOptions["private_key_file"] = key_file sfOptions["private_key_file_pwd"] = key_passphrase sfOptions["db"] ="**_DB" sfOptions["warehouse"] ="****_WHS" sfOptions["schema"] ="***_SHR" sfOptions["role"] ="**_ROLE" sfOptions["numPartitions"]="10" sfOptions["partitionColumn"] = "***_TRANS_ID" sfOptions["lowerBound"] = lowerbound sfOptions["upperBound"] = upperbound # 替换原来的query参数,需要自定义过滤逻辑时可以写子查询并加别名 sfOptions["dbtable"] = "***_shr.SPRK_TST" df = spark.read.format('jdbc') \ .options(**sfOptions) \ .load()
方案二:切换为Snowflake官方Spark连接器(更稳定,推荐)
Snowflake官方Spark连接器原生支持分区参数,兼容性和性能优于原生JDBC对接方式,修改后代码示例:
sfOptions = { "sfURL" : "**************.privatelink.snowflakecomputing.com", "sfUser" : "**01d", "sfPrivateKeyFile": key_file, "sfPrivateKeyFilePwd": key_passphrase, "sfDatabase" : "**_DB", "sfWarehouse" : "****_WHS", "sfSchema" : "***_SHR", "sfRole" : "**_ROLE", "dbtable" : "SPRK_TST", "partition_column" : "***_TRANS_ID", "lower_bound" : lowerbound, "upper_bound" : upperbound, "numPartitions": "10" } df = spark.read.format("net.snowflake.spark.snowflake") \ .options(**sfOptions) \ .load()
注意:分区列请选择数值类型、值分布均匀的字段,否则会出现分区数据量差异过大的问题,影响执行效率。
内容的提问来源于stack exchange,提问作者gopinath kolanchi
相关产品推荐
相关产品推荐

