【发布时间】:2015-09-29 06:18:57
【问题描述】:
我有一个非常简单的 SparkSQL 设置连接到 Postgres 数据库,我试图从表中获取一个 DataFrame,该 DataFrame 具有多个 X 分区(假设为 2)。代码如下:
Map<String, String> options = new HashMap<String, String>();
options.put("url", DB_URL);
options.put("driver", POSTGRES_DRIVER);
options.put("dbtable", "select ID, OTHER from TABLE limit 1000");
options.put("partitionColumn", "ID");
options.put("lowerBound", "100");
options.put("upperBound", "500");
options.put("numPartitions","2");
DataFrame housingDataFrame = sqlContext.read().format("jdbc").options(options).load();
由于某种原因,DataFrame 的一个分区几乎包含所有行。
据我所知lowerBound/upperBound 是用于微调的参数。在 SparkSQL 的文档(Spark 1.4.0 - spark-sql_2.11)中,它说它们用于定义步幅,而不是过滤/范围分区列。但这提出了几个问题:
- 步幅是 Spark 为每个执行程序(分区)查询数据库的频率(每个查询返回的元素数)?
- 如果不是,此参数的用途是什么,它们依赖于什么以及如何以稳定的方式平衡我的 DataFrame 分区(不要求所有分区包含相同数量的元素,只是存在一个平衡 -例如 2 个分区 100 个元素 55/45 、 60/40 甚至 65/35 都可以)
似乎无法找到这些问题的明确答案,想知道你们中的一些人是否可以为我澄清这一点,因为现在在处理 X 百万行时影响我的集群性能,所有繁重的工作都过去了给一个执行者。
干杯并感谢您的宝贵时间。
【问题讨论】:
标签: postgresql apache-spark apache-spark-sql partition