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的失败输出标签,导致异常无法被下游收集。改造逻辑如下:
- 定义输出标签,标记ES写入的成功、失败分支
// 定义ES写入成功、失败集合的标签 final TupleTag<String> ES_WRITE_SUCCESS = new TupleTag<String>(){}; final TupleTag<FailsafeElement<String, String>> ES_WRITE_FAILURE = new TupleTag<FailsafeElement<String, String>>(){};
- 改造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);
- 编写错误分类工具方法,对捕获的异常打类型标
/** * 划分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"; }
- 格式化错误记录,分别适配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()); }) );
- 将格式化后的错误日志分别写入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
相关产品推荐
相关产品推荐

