You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Apache Beam中启用ElasticSearchIO并行读取?

ElasticsearchIO读取并行扩展优化问题

我有一个可在Google Dataflow和DirectRunner中运行的简单管道,用于与ElasticSearch集群交互处理数据,基本流程如下:

  • 使用ElasticIO连接器从ElasticSearch读取文档
  • 将文档序列化为内部Protobuf格式
  • 执行无外部依赖的转换,输出另一种Protobuf格式
  • 将最终Protobuf写入另一个ElasticSearch索引

该管道在测试场景运行正常,但扩展至处理数亿级文档时,并行度无法提升。在Google Dataflow中运行超过5小时,处理速率仅为每秒50条,效率极低。我们内部单实例系统可达到每秒1万-2万条的处理速率,且其他应用从ElasticSearch读取数据速度很快,暂不认为ElasticSearch集群是瓶颈。

已尝试的优化操作

  • 增加numWorkers:临时提升了worker数量,但因ElasticSearch读取无法提供足够数据维持worker负载,worker数量又降回原值
  • 修改ElasticSearch读取的批量大小:无任何效果

管道配置代码

PCollection<String> dataCollection = pipeline.apply("Reading From ElasticSearch", ElasticsearchIO.read()
        .withConnectionConfiguration(esReadConnection)
        .withBatchSize(options.getBatchSize())
        .withScrollKeepalive(scrollTime)
        .withQuery(options.getQuery())//query needs to be stringified json including with the "query" element
        .withMetadata()
    );
dataCollection.apply("Serialize", ParDo.of(new JsonToProto<>(searchHitTag, failTag, SearchHit::newBuilder)));

问题根源排查

查看ElasticsearchIO的拆分逻辑后发现,问题出在配置使用的索引名称是别名(部分场景为数据流)。getEstimatedSizeBytes方法中的以下代码未能覆盖这类场景:

JsonNode indexStats =
          statsJson.path("indices").path(connectionConfiguration.getIndex()).path("primaries");
long indexSize = indexStats.path("store").path("size_in_bytes").asLong();

对集群执行统计查询后,indices键下返回的是实际索引列表,但其中没有与配置中完全匹配的名称。该方法直接用配置的索引名获取统计数据的逻辑存在缺陷,因为在ElasticSearch的合理使用场景中,查询用的索引名称/模式常与实际索引名称不一致(比如别名、数据流),这会导致无法正确估算数据量,进而无法拆分出足够的分片来提升并行度,数据流场景也会因此无法获得高性能。

可行修复方案

  • 改用_all.primaries路径获取整体统计数据,indexSize路径保持不变
  • 遍历indices对象中的所有索引,累加每个索引的size_in_bytes

目前已针对Apache Beam的ElasticSearchIO连接器提交修复工单,并附上了可解决此问题的补丁。

内容的提问来源于stack exchange,提问作者Jeff

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.13 04:40:33