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});
     }

Reply via email to