This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12312-33f430649724995f08c2ea876b557f7e6d0646c6 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 065a3df4b39dc116c3d432b15d5c7c4eadec120c Author: Jast <[email protected]> AuthorDate: Fri Sep 18 14:41:00 2026 +0000 [Fix][Connector-V2] Propagate HTTP sink write failures (#12312) Co-authored-by: zhangshenghang <[email protected]> --- docs/en/connectors/sink/Http.md | 2 ++ .../introduction/concepts/incompatible-changes.md | 6 ++++ docs/zh/connectors/sink/Http.md | 2 ++ .../introduction/concepts/incompatible-changes.md | 6 ++++ .../seatunnel/http/sink/HttpSinkWriter.java | 28 +++++++++------ .../http/sink/HttpSinkBatchWriterTest.java | 40 ++++++++++++++++++++++ 6 files changed, 74 insertions(+), 10 deletions(-) diff --git a/docs/en/connectors/sink/Http.md b/docs/en/connectors/sink/Http.md index b06b2be3d0..3e83b2fc63 100644 --- a/docs/en/connectors/sink/Http.md +++ b/docs/en/connectors/sink/Http.md @@ -54,6 +54,8 @@ They can be downloaded via install-plugin.sh or from the Maven central repositor The Http sink always sends `POST` requests. Each upstream row is converted to JSON and used as the request body. When `array_mode = true`, rows are accumulated into a JSON array before sending; `batch_size` controls the maximum number of rows in one request. +The sink treats a non-200 HTTP response or a request exception as a write failure and reports it to the engine. A failed batch is not acknowledged as successfully flushed. The connector does not provide exactly-once delivery, so receivers should remain idempotent. + simple: ```hocon diff --git a/docs/en/introduction/concepts/incompatible-changes.md b/docs/en/introduction/concepts/incompatible-changes.md index ebe8cebd76..f2fedd0aca 100644 --- a/docs/en/introduction/concepts/incompatible-changes.md +++ b/docs/en/introduction/concepts/incompatible-changes.md @@ -146,6 +146,12 @@ You need to check this document before you upgrade to related version. ### Connector Changes +- **Behavior change: HTTP sink write failures now fail the task instead of being silently dropped** + - **Affected component**: `seatunnel-connectors-v2/connector-http/connector-http-base` + - **Description**: Previously, `HttpSinkWriter.doHttpRequest` handled both a non-200 HTTP response and any request exception (network error, timeout, serialization error) by logging at `error` level and returning normally, so the failed row/batch was silently dropped while the job kept running and checkpoints completed. The writer now throws `HttpConnectorException` (`REQUEST_FAILED`) for both cases, so the failure propagates to the engine and fails the task/job. + - **Impact**: Jobs whose downstream HTTP endpoint occasionally returns non-200 or is occasionally unreachable used to keep running with silent data loss; after this change they fail loudly at the first failed write. The connector still has no built-in retry or dead-letter mechanism, so re-submitting a failed job may deliver rows that succeeded before the failure again — make sure the receiver tolerates duplicate delivery on retry/restart. + - **Migration Guide**: No configuration change is required. If your endpoint is expected to return non-200 responses as part of normal operation, handle them upstream of the sink or add an external retry mechanism before upgrading. + - **Breaking Change: BigQuery Sink Connector — default schema save mode introduces automatic table creation** - **Affected component**: `seatunnel-connectors-v2/connector-bigquery` - **Description**: The BigQuery sink connector (`connector-bigquery`) now implements `SupportSaveMode` with support for `schema_save_mode` and `data_save_mode`. The default `schema_save_mode` is set to `CREATE_SCHEMA_WHEN_NOT_EXIST`. diff --git a/docs/zh/connectors/sink/Http.md b/docs/zh/connectors/sink/Http.md index afb3be4d1c..e18718a15b 100644 --- a/docs/zh/connectors/sink/Http.md +++ b/docs/zh/connectors/sink/Http.md @@ -53,6 +53,8 @@ import ChangeLog from '../changelog/connector-http.md'; Http Sink 固定发送 `POST` 请求。每条上游数据会被转换成 JSON 作为请求体;当 `array_mode = true` 时,会先把多条数据攒成 JSON 数组再发送,`batch_size` 控制单次请求最多包含多少条数据。 +当 HTTP 响应不是 200,或请求过程中发生异常时,Sink 会将写入失败上报给引擎;失败批次不会被当作已成功发送。该连接器不提供精确一次投递保证,接收端应保证幂等。 + 简单示例: ```hocon diff --git a/docs/zh/introduction/concepts/incompatible-changes.md b/docs/zh/introduction/concepts/incompatible-changes.md index ea8fc9d80b..bec69ded47 100644 --- a/docs/zh/introduction/concepts/incompatible-changes.md +++ b/docs/zh/introduction/concepts/incompatible-changes.md @@ -129,6 +129,12 @@ ### 连接器变更 +- **行为变更:HTTP Sink 写入失败现在会使任务失败,而不再被静默丢弃** + - **影响范围**:`seatunnel-connectors-v2/connector-http/connector-http-base` + - **变更说明**:此前 `HttpSinkWriter.doHttpRequest` 对非 200 的 HTTP 响应和任何请求异常(网络错误、超时、序列化错误)都只记录 `error` 日志后正常返回,导致失败的行/批次被静默丢弃,而作业继续运行、checkpoint 正常完成。现在这两种情况都会抛出 `HttpConnectorException`(`REQUEST_FAILED`),失败会传播到引擎并使任务/作业失败。 + - **影响**:下游 HTTP 端点偶发返回非 200 或偶发不可达的作业,此前会带着静默丢数据继续运行;升级后会在第一次写入失败时大声失败。该连接器仍没有内置重试或死信机制,重新提交失败的作业可能重复投递失败前已成功的行——请确保接收端能够容忍重试/重启时的重复投递。 + - **迁移指南**:无需更改配置。如果您的端点在正常业务中就会返回非 200 响应,请在 Sink 之前的环节处理这些响应,或在升级前引入外部重试机制。 + - **破坏性变更:ORC 文件 Sink 保留嵌套 Struct 字段名的大小写** - **影响范围**:`seatunnel-connectors-v2/connector-file/connector-file-base`(所有共享 `OrcWriteStrategy` 的 File/HDFS/S3/OSS ORC Sink) - **变更说明**:此前,`OrcWriteStrategy.buildFieldWithRowType(...)` 在构建 ORC Schema 时,会将每个嵌套 `ROW`(struct)字段名强制转为小写,因此声明为 `MD5` 的嵌套字段在文件 footer 中被持久化为 `md5`。下游消费者按原始大小写名称读取该列时会得到 null/缺失值。本次移除了递归嵌套字段分支上的 `.toLowerCase()` 调用,嵌套 struct 字段名将按原始大小写写入文件 Schema。 diff --git a/seatunnel-connectors-v2/connector-http/connector-http-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/http/sink/HttpSinkWriter.java b/seatunnel-connectors-v2/connector-http/connector-http-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/http/sink/HttpSinkWriter.java index cde7eabe36..f9b81328aa 100644 --- a/seatunnel-connectors-v2/connector-http/connector-http-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/http/sink/HttpSinkWriter.java +++ b/seatunnel-connectors-v2/connector-http/connector-http-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/http/sink/HttpSinkWriter.java @@ -29,6 +29,8 @@ import org.apache.seatunnel.connectors.seatunnel.common.sink.AbstractSinkWriter; import org.apache.seatunnel.connectors.seatunnel.http.client.HttpClientProvider; import org.apache.seatunnel.connectors.seatunnel.http.client.HttpResponse; import org.apache.seatunnel.connectors.seatunnel.http.config.HttpParameter; +import org.apache.seatunnel.connectors.seatunnel.http.exception.HttpConnectorErrorCode; +import org.apache.seatunnel.connectors.seatunnel.http.exception.HttpConnectorException; import org.apache.seatunnel.format.json.JsonSerializationSchema; import lombok.extern.slf4j.Slf4j; @@ -142,22 +144,28 @@ public class HttpSinkWriter extends AbstractSinkWriter<SeaTunnelRow, Void> if (HttpResponse.STATUS_OK == response.getCode()) { return; } - log.error( - "http client execute exception, http response status code:[{}], content:[{}]", - response.getCode(), - response.getContent()); + String message = + String.format( + "http client execute exception, http response status code:[%s], content:[%s]", + response.getCode(), response.getContent()); + throw new HttpConnectorException(HttpConnectorErrorCode.REQUEST_FAILED, message); + } catch (HttpConnectorException e) { + throw e; } catch (Exception e) { - log.error(e.getMessage(), e); + throw new HttpConnectorException(HttpConnectorErrorCode.REQUEST_FAILED, e); } } @Override public void close() throws IOException { - if (arrayMode) { - flush(); - } - if (Objects.nonNull(httpClient)) { - httpClient.close(); + try { + if (arrayMode) { + flush(); + } + } finally { + if (Objects.nonNull(httpClient)) { + httpClient.close(); + } } } diff --git a/seatunnel-connectors-v2/connector-http/connector-http-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/http/sink/HttpSinkBatchWriterTest.java b/seatunnel-connectors-v2/connector-http/connector-http-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/http/sink/HttpSinkBatchWriterTest.java index ee8769f595..d1145327c5 100644 --- a/seatunnel-connectors-v2/connector-http/connector-http-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/http/sink/HttpSinkBatchWriterTest.java +++ b/seatunnel-connectors-v2/connector-http/connector-http-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/http/sink/HttpSinkBatchWriterTest.java @@ -24,6 +24,7 @@ import org.apache.seatunnel.api.table.type.SeaTunnelRowType; import org.apache.seatunnel.connectors.seatunnel.http.client.HttpClientProvider; import org.apache.seatunnel.connectors.seatunnel.http.client.HttpResponse; import org.apache.seatunnel.connectors.seatunnel.http.config.HttpParameter; +import org.apache.seatunnel.connectors.seatunnel.http.exception.HttpConnectorException; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -36,11 +37,13 @@ import org.mockito.junit.jupiter.MockitoExtension; import org.mockito.junit.jupiter.MockitoSettings; import org.mockito.quality.Strictness; +import java.io.IOException; import java.util.HashMap; import java.util.Map; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; @@ -183,6 +186,43 @@ public class HttpSinkBatchWriterTest { assertTrue(requestBody.endsWith("]")); } + @Test + public void testArrayModeFailureIsPropagatedDuringCheckpoint() throws Exception { + httpParameter.setArrayMode(true); + httpParameter.setBatchSize(BATCH_SIZE); + sinkWriter = new TestableHttpSinkWriter(rowType, httpParameter); + HttpResponse failedResponse = new HttpResponse(500, "server error"); + when(httpClientProvider.doPost(anyString(), any(), anyString())).thenReturn(failedResponse); + + sinkWriter.write(createTestRow(1, "user1", 20)); + + HttpConnectorException exception = + assertThrows(HttpConnectorException.class, () -> sinkWriter.prepareCommit()); + + assertTrue(exception.getMessage().contains("http response status code:[500]")); + HttpResponse successfulResponse = new HttpResponse(HttpResponse.STATUS_OK); + when(httpClientProvider.doPost(anyString(), any(), anyString())) + .thenReturn(successfulResponse); + + sinkWriter.prepareCommit(); + + verify(httpClientProvider, times(2)).doPost(eq(TEST_URL), any(), anyString()); + } + + @Test + public void testObjectModeRequestExceptionIsPropagated() throws Exception { + sinkWriter = new TestableHttpSinkWriter(rowType, httpParameter); + when(httpClientProvider.doPost(anyString(), any(), anyString())) + .thenThrow(new IOException("connection refused")); + + HttpConnectorException exception = + assertThrows( + HttpConnectorException.class, + () -> sinkWriter.write(createTestRow(1, "user1", 20))); + + assertEquals("connection refused", exception.getCause().getMessage()); + } + private SeaTunnelRow createTestRow(int id, String name, int age) { return new SeaTunnelRow(new Object[] {id, name, age}); }
