【发布时间】:2018-05-13 21:22:56
【问题描述】:
我使用 spark 读取 elasticsearch.Like
select col from index limit 10;
问题是索引非常大,它包含1000亿行。并且spark生成数千个任务来完成这项工作。
我只需要 10 行,即使 1 个任务返回 10 行也可以完成工作。我不需要这么多任务。
即使是限制1,限制也很慢。
代码:
sql = select col from index limit 10
sqlExecListener.sparkSession.sql(sql).createOrReplaceTempView(tempTable)
【问题讨论】:
-
您是否尝试过明确设置分区大小?
-
@JustinPihony 是的,我设置了 es_input_max_docs_per_partition=5000,看来 total_rows_es_contains = es_input_max_docs_per_partition * num_of_partitions
-
如果你使用 push.down=True 确保 double.filtering=False 因为它可能会阻止限制被推低。通过调用 df.explain(True) 检查你的物理计划,并确保弹性搜索和限制之间没有过滤
标签: apache-spark elasticsearch apache-spark-sql spark-submit