【问题标题】:Why doesn't Spark Mongo connector push down filters?为什么 Spark Mongo 连接器不下推过滤器?
【发布时间】:2019-04-18 13:56:19
【问题描述】:

我有一个大型 Mongo 集合,我想在我的 Spark 应用程序中使用它,使用 Spark Mongo 连接器。该集合非常大(> 10 GB)并且有每日数据,在original_item.CreatedDate字段上有一个索引。在 Mongo 中选择几天的查询非常快(不到一秒)。但是,当我使用数据帧编写相同的查询时,该过滤器不会下推到 Mongo,导致性能极慢,因为 Spark 显然会获取整个集合并自行过滤。

查询如下所示:

collection
      .filter("original_item.CreatedDate  > %s" % str(start_date_timestamp_ms)) \
      .filter("original_item.CreatedDate  < %s" % str(end_date_timestamp_ms)) \
      .select(...)

在物理计划中,我看到: PushedFilters: [IsNotNull(original_item)]

当我对该集合的另一个字段进行过滤并进行类似查询时,mongo 成功将其下推 - PushedFilters: [IsNotNull(original_item), IsNotNull(doc_type), EqualTo(doc_type,case)]!

会不会是 GreaterThan 过滤器推送不受 Mongo Spark 连接器支持,或者它存在错误?

谢谢!

【问题讨论】:

    标签: database mongodb apache-spark


    【解决方案1】:

    导致您的问题的不是GreaterThan,而是过滤器位于嵌套字段上。 doc_type 上的过滤器有效,因为它没有嵌套。这显然是 Spark 中的 Catalyst 引擎的问题,而不是 Mongo 连接器的问题。它也会影响例如 Parquet 中的谓词下推。

    有关更多详细信息,请参阅 Spark Jira 中的以下讨论。

    Spark 19638

    Spark 17636

    【讨论】:

      猜你喜欢
      • 2016-06-23
      • 1970-01-01
      • 1970-01-01
      • 2019-06-02
      • 1970-01-01
      • 1970-01-01
      • 2023-03-14
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多