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

Java Dataflow/Beam实现ElasticSearch错误重定向至BigQuery/错误主题

GCP Dataflow 对接ElasticSearch错误捕获分流实现

需求说明

需基于GCP平台的Dataflow/Apache Beam框架,实现ElasticSearch操作错误的全量捕获与分流能力:

  • 覆盖全量操作异常,包含可恢复、不可恢复两类,典型场景如ConnectTimeOut、Keystore加载错误等
  • 捕获的错误信息需写入BigQuery表留存全量日志,同时投递至指定错误主题,避免数据丢失
  • 明确ElasticSearch错误分类处理逻辑,以及错误重定向至BigQuery完成日志落盘的具体实现方法

现有实现

当前已完成自定义流水线模板开发,实现从PubSub读取数据写入ElasticSearch的核心流程,代码片段如下:

/*
* Step #1: Read from a PubSub subscription.
*/
PCollection messages = null;

if (options.getUseSubscription()) {

  messages =

          pipeline.apply(

                  "ReadPubSubSubscription",

                  PubsubIO.readMessagesWithAttributes()
                          .fromSubscription(options.getInputSubscription()));

}

/*
* Step #2: Transform the PubsubMessages into Json documents.
*/
PCollectionTuple convertedPubsubMessages = messages.apply(

      "ConvertMessageToJsonDocument",new PubSubMessageToJsonDocument(options));

/*
* Step #3a: Write Json documents into Elasticsearch using {@link ElasticsearchTransforms.WriteToElasticsearch}.
*/
convertedPubsubMessages

      .get(TRANSFORM_OUT)

      .apply(
              "GetJsonDocuments",
              MapElements.into(TypeDescriptors.strings()).via(FailsafeElement::getPayload)
)
.apply(

              "WriteToElasticsearch",

              ElasticsearchTransforms.WriteToElasticsearch.newBuilder()
                      .setOptions(options.as(WriteToElasticsearchOptions.class))

                      .build());

实现方案

错误分类规则

首先明确两类异常的划分标准,方便后续做差异化处理:

  • 可恢复异常:属于瞬时故障,包含连接超时、集群节点临时不可用、请求限流(429状态码)、临时网络抖动等,这类异常框架默认会触发重试,超过重试阈值后进入失败分支
  • 不可恢复异常:属于永久故障,包含Keystore/证书加载失败、索引不存在、字段映射不匹配、权限校验失败、请求格式非法等,这类异常重试无意义,直接进入错误处理分支

代码改造步骤

原有实现的核心问题是在写入ES前通过MapElements单独提取了payload,丢弃了FailsafeElement携带的错误上下文,且没有配置ES写入Transform的失败输出标签,导致异常无法被下游收集。改造逻辑如下:

  1. 定义输出标签,标记ES写入的成功、失败分支
// 定义ES写入成功、失败集合的标签
final TupleTag<String> ES_WRITE_SUCCESS = new TupleTag<String>(){};
final TupleTag<FailsafeElement<String, String>> ES_WRITE_FAILURE = new TupleTag<FailsafeElement<String, String>>(){};
  1. 改造ES写入逻辑,保留错误上下文,开启失败记录收集
/*
* Step #3: 写入Elasticsearch并收集全量失败记录
*/
PCollectionTuple esWriteResults = convertedPubsubMessages
      .get(TRANSFORM_OUT)
      .apply(
              "WriteToElasticsearchWithErrorCapture",
              ElasticsearchTransforms.WriteToElasticsearch.newBuilder()
                      .setOptions(options.as(WriteToElasticsearchOptions.class))
                      // 绑定成功、失败输出标签
                      .setFailureTag(ES_WRITE_FAILURE)
                      .setSuccessTag(ES_WRITE_SUCCESS)
                      .build()
      );

// 提取所有ES写入失败的记录
PCollection<FailsafeElement<String, String>> esFailedRecords = esWriteResults.get(ES_WRITE_FAILURE);
  1. 编写错误分类工具方法,对捕获的异常打类型标
/**
 * 划分ES操作异常类型
 */
private String classifyEsError(Throwable e) {
    if (e instanceof ConnectTimeoutException || e instanceof SocketTimeoutException 
        || e instanceof EsRejectedExecutionException) {
        return "RECOVERABLE";
    }
    if (e instanceof KeyStoreException || (e instanceof ElasticsearchStatusException 
        && ((ElasticsearchStatusException)e).status().getStatus() >= 400 
        && ((ElasticsearchStatusException)e).status().getStatus() != 429)) {
        return "NON_RECOVERABLE";
    }
    return "UNKNOWN";
}
  1. 格式化错误记录,分别适配BigQuery、PubSub的写入格式
/*
* Step #4: 格式化错误记录适配下游存储
*/
// 转换为BigQuery表行格式
PCollection<TableRow> errorLogsForBQ = esFailedRecords
      .apply(
              "FormatErrorToBigQueryRow",
              MapElements.into(TypeDescriptors.rows())
                      .via(failedElement -> {
                          String errorMsg = failedElement.getErrorMessage();
                          Throwable err = failedElement.getError();
                          String stackTrace = Throwables.getStackTraceAsString(err);
                          String originalPayload = failedElement.getPayload();
                          Long eventTime = Instant.now().getMillis() / 1000;
                          String errorType = classifyEsError(err);
                          
                          return new TableRow()
                                  .set("raw_payload", originalPayload)
                                  .set("error_message", errorMsg)
                                  .set("error_stack", stackTrace)
                                  .set("error_type", errorType)
                                  .set("event_timestamp", eventTime);
                      })
      );

// 转换为Pub/Sub消息格式
PCollection<PubsubMessage> errorMessagesForPubsub = esFailedRecords
      .apply(
              "FormatErrorToPubsubMessage",
              MapElements.into(TypeDescriptors.pubsubMessages())
                      .via(failedElement -> {
                          String errContent = String.format("{\"raw_payload\":\"%s\",\"error_msg\":\"%s\",\"error_type\":\"%s\"}",
                                  failedElement.getPayload(),
                                  failedElement.getErrorMessage(),
                                  classifyEsError(failedElement.getError()));
                          return new PubsubMessage(errContent.getBytes(StandardCharsets.UTF_8), ImmutableMap.of());
                      })
      );
  1. 将格式化后的错误日志分别写入BigQuery、错误Pub/Sub主题
/*
* Step #5: 错误日志多目的地落地
*/
// 写入BigQuery留存全量日志
errorLogsForBQ.apply(
        "WriteErrorToBigQuery",
        BigQueryIO.writeTableRows()
                .to(options.getErrorBigQueryTable())
                .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND)
                .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED)
                .withSchema(getErrorTableSchema()) // 提前定义与组装字段匹配的BQ表schema
);

// 投递至错误Pub/Sub主题
errorMessagesForPubsub.apply(
        "PublishErrorToPubsubTopic",
        PubsubIO.writeMessages().to(options.getErrorPubsubTopic())
);

注意事项

  • 不要在写入ES前剥离FailsafeElement的上下文单独提取payload,否则会丢失框架携带的原始消息、异常栈、错误信息,导致后续无法定位问题根因
  • 可根据业务需要,从失败记录中单独过滤RECOVERABLE类型的异常接入死信队列做延迟重试,NON_RECOVERABLE类型异常可直接对接告警规则
  • 配置ES写入Transform时建议设置合理的最大重试次数、退避间隔,避免瞬时故障直接进入错误分支

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 09:54:25