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-12286-08a38f1b8c5353228d3e1ba11e2eebb20619844f
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit 62487f093e7b73393d3e5e08282ad3be6b46fc54
Author: Goutam Adwant <[email protected]>
AuthorDate: Sat Sep 19 00:51:21 2026 +0000

    [Fix][Connector-V2] Parse offset timestamps in Iceberg delete predicates 
(#12286)
    
    Signed-off-by: Goutam Adwant <[email protected]>
---
 docs/en/connectors/sink/Iceberg.md                 |   4 +
 docs/zh/connectors/sink/Iceberg.md                 |   4 +
 .../seatunnel/iceberg/utils/ExpressionUtils.java   |  17 +++
 .../catalog/IcebergTimestampDeleteTest.java        | 161 +++++++++++++++++++++
 .../iceberg/utils/ExpressionUtilsTest.java         | 102 +++++++++++++
 5 files changed, 288 insertions(+)

diff --git a/docs/en/connectors/sink/Iceberg.md 
b/docs/en/connectors/sink/Iceberg.md
index 26d3395b58..8c4be05d2e 100644
--- a/docs/en/connectors/sink/Iceberg.md
+++ b/docs/en/connectors/sink/Iceberg.md
@@ -104,6 +104,10 @@ Write commits to the specified Iceberg branch. Leave it 
empty to commit to the t
 
 When `data_save_mode = CUSTOM_PROCESSING`, configure the `delete` SQL that 
removes the target data before the sink writes. This option is required in that 
mode.
 
+For Iceberg `timestamptz` columns, comparison and `IN`/`NOT IN` predicates 
accept timestamp literals with an offset, for example `event_time >= 
'2026-09-12 10:00:00.123456+05:30'`. Offsets are converted to the corresponding 
UTC instant at microsecond precision. Offset-free literals retain their 
existing UTC interpretation. Iceberg `timestamp` columns without a time zone 
continue to require offset-free literals.
+
+The existing file-level delete restriction still applies: Iceberg rejects a 
delete when only some rows in a data file match the condition.
+
 ### krb5_path [string]
 
 The path of `krb5.conf`, used for Kerberos authentication.
diff --git a/docs/zh/connectors/sink/Iceberg.md 
b/docs/zh/connectors/sink/Iceberg.md
index ad3e7d35fb..13d5bdd6ff 100644
--- a/docs/zh/connectors/sink/Iceberg.md
+++ b/docs/zh/connectors/sink/Iceberg.md
@@ -104,6 +104,10 @@ libfb303-xxx.jar
 
 当 `data_save_mode = CUSTOM_PROCESSING` 时,配置在 Sink 写入前删除目标数据的 `delete` 
SQL。该模式下必须配置此选项。
 
+对于 Iceberg `timestamptz` 列,比较以及 `IN`/`NOT IN` 条件支持带偏移量的时间戳字面量,例如 `event_time 
>= '2026-09-12 10:00:00.123456+05:30'`。偏移量会转换为对应的 UTC 时间点,精度为微秒。不带偏移量的字面量保持现有的 
UTC 解释方式。不带时区的 Iceberg `timestamp` 列仍要求使用不带偏移量的字面量。
+
+现有的文件级删除限制仍然适用:如果数据文件中只有部分行匹配条件,Iceberg 会拒绝删除。
+
 ### krb5_path [string]
 
 `krb5.conf` 文件的路径,用于 Kerberos 认证。
diff --git 
a/seatunnel-connectors-v2/connector-iceberg/src/main/java/org/apache/seatunnel/connectors/seatunnel/iceberg/utils/ExpressionUtils.java
 
b/seatunnel-connectors-v2/connector-iceberg/src/main/java/org/apache/seatunnel/connectors/seatunnel/iceberg/utils/ExpressionUtils.java
index a5de9c2b4a..606bcdd0f4 100644
--- 
a/seatunnel-connectors-v2/connector-iceberg/src/main/java/org/apache/seatunnel/connectors/seatunnel/iceberg/utils/ExpressionUtils.java
+++ 
b/seatunnel-connectors-v2/connector-iceberg/src/main/java/org/apache/seatunnel/connectors/seatunnel/iceberg/utils/ExpressionUtils.java
@@ -55,8 +55,10 @@ import java.math.BigDecimal;
 import java.time.LocalDate;
 import java.time.LocalDateTime;
 import java.time.LocalTime;
+import java.time.OffsetDateTime;
 import java.time.format.DateTimeFormatter;
 import java.time.format.DateTimeFormatterBuilder;
+import java.time.temporal.ChronoField;
 import java.util.ArrayList;
 import java.util.List;
 import java.util.stream.Collectors;
@@ -73,6 +75,15 @@ public class ExpressionUtils {
                     .append(ISO_LOCAL_TIME)
                     .toFormatter();
 
+    private static final DateTimeFormatter OFFSET_DATE_TIME_FORMATTER =
+            new DateTimeFormatterBuilder()
+                    .append(LOCAL_DATE_TIME_FORMATTER)
+                    .optionalStart()
+                    .appendOffsetId()
+                    .optionalEnd()
+                    .parseDefaulting(ChronoField.OFFSET_SECONDS, 0)
+                    .toFormatter();
+
     public static List<String> parseSelectColumns(String selectQuery) {
         if (StringUtils.isNotBlank(selectQuery)) {
             try {
@@ -274,6 +285,12 @@ public class ExpressionUtils {
                 }
             case TIMESTAMP:
                 if (valueExpression instanceof StringValue) {
+                    if (((Types.TimestampType) 
icebergColumn.type()).shouldAdjustToUTC()) {
+                        return DateTimeUtil.microsFromTimestamptz(
+                                OffsetDateTime.parse(
+                                        ((StringValue) 
valueExpression).getValue(),
+                                        OFFSET_DATE_TIME_FORMATTER));
+                    }
                     LocalDateTime dateTime =
                             LocalDateTime.parse(
                                     ((StringValue) valueExpression).getValue(),
diff --git 
a/seatunnel-connectors-v2/connector-iceberg/src/test/java/org/apache/seatunnel/connectors/seatunnel/iceberg/catalog/IcebergTimestampDeleteTest.java
 
b/seatunnel-connectors-v2/connector-iceberg/src/test/java/org/apache/seatunnel/connectors/seatunnel/iceberg/catalog/IcebergTimestampDeleteTest.java
new file mode 100644
index 0000000000..a77281a1ac
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-iceberg/src/test/java/org/apache/seatunnel/connectors/seatunnel/iceberg/catalog/IcebergTimestampDeleteTest.java
@@ -0,0 +1,161 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.iceberg.catalog;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.table.catalog.TablePath;
+import org.apache.seatunnel.connectors.seatunnel.iceberg.IcebergCatalogLoader;
+import 
org.apache.seatunnel.connectors.seatunnel.iceberg.config.IcebergCommonOptions;
+
+import org.apache.iceberg.FileFormat;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.catalog.Namespace;
+import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.data.GenericAppenderFactory;
+import org.apache.iceberg.data.GenericRecord;
+import org.apache.iceberg.data.IcebergGenerics;
+import org.apache.iceberg.data.Record;
+import org.apache.iceberg.encryption.EncryptedFiles;
+import org.apache.iceberg.exceptions.ValidationException;
+import org.apache.iceberg.inmemory.InMemoryCatalog;
+import org.apache.iceberg.io.CloseableIterable;
+import org.apache.iceberg.io.DataWriter;
+import org.apache.iceberg.types.Types;
+
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedConstruction;
+import org.mockito.Mockito;
+
+import java.time.OffsetDateTime;
+import java.time.format.DateTimeParseException;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.Map;
+import java.util.Set;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+class IcebergTimestampDeleteTest {
+    @Test
+    void deleteOffsetTimestampWindowPreservesOtherRecords() throws Exception {
+        Map<String, Object> config = new HashMap<>();
+        config.put(IcebergCommonOptions.KEY_CATALOG_NAME.key(), "test");
+        config.put(IcebergCommonOptions.CATALOG_PROPS.key(), 
Collections.emptyMap());
+        // Keep real Iceberg commits and Parquet I/O without Hadoop's native 
filesystem tools.
+        try (InMemoryCatalog setup = new InMemoryCatalog();
+                MockedConstruction<IcebergCatalogLoader> ignored =
+                        Mockito.mockConstruction(
+                                IcebergCatalogLoader.class,
+                                (loader, context) ->
+                                        
Mockito.when(loader.loadCatalog()).thenReturn(setup))) {
+            setup.initialize("test", Collections.emptyMap());
+            setup.createNamespace(Namespace.of("test"));
+            IcebergCatalog catalog = new IcebergCatalog("test", 
ReadonlyConfig.fromMap(config));
+            Schema schema =
+                    new Schema(
+                            Types.NestedField.required(1, "id", 
Types.IntegerType.get()),
+                            Types.NestedField.optional(
+                                    2, "event_time", 
Types.TimestampType.withZone()));
+            Table table =
+                    setup.createTable(
+                            TableIdentifier.of("test", "events"),
+                            schema,
+                            PartitionSpec.unpartitioned());
+            OffsetDateTime boundary = 
OffsetDateTime.parse("2026-09-12T04:30:00.123456Z");
+            appendRecord(table, 1, boundary.minusNanos(1000));
+            appendRecord(table, 2, boundary);
+            appendRecord(table, 3, boundary.plusNanos(1000));
+            appendRecord(table, 4, null);
+            catalog.open();
+            try {
+                TablePath path = TablePath.of("test", "events");
+                long snapshotId = table.currentSnapshot().snapshotId();
+                assertThrows(
+                        DateTimeParseException.class,
+                        () ->
+                                catalog.executeSql(
+                                        path,
+                                        "DELETE FROM test.events WHERE 
event_time >= 'invalid'"));
+                table.refresh();
+                assertEquals(snapshotId, table.currentSnapshot().snapshotId());
+                assertEquals(new HashSet<>(Arrays.asList(1, 2, 3, 4)), 
readIds(table));
+                catalog.executeSql(
+                        path,
+                        "DELETE FROM test.events WHERE event_time >= 
'2026-09-12 10:00:00.123456+05:30' AND event_time < '2026-09-11 
23:30:00.123457-05:00'");
+                table.refresh();
+                assertEquals(new HashSet<>(Arrays.asList(1, 3, 4)), 
readIds(table));
+                appendRecords(table, 5, new OffsetDateTime[] {boundary, 
boundary.plusNanos(1000)});
+                long mixedFileSnapshotId = 
table.currentSnapshot().snapshotId();
+                assertThrows(
+                        ValidationException.class,
+                        () ->
+                                catalog.executeSql(
+                                        path,
+                                        "DELETE FROM test.events WHERE 
event_time = '2026-09-12 10:00:00.123456+05:30'"));
+                table.refresh();
+                assertEquals(mixedFileSnapshotId, 
table.currentSnapshot().snapshotId());
+                assertEquals(new HashSet<>(Arrays.asList(1, 3, 4, 5, 6)), 
readIds(table));
+            } finally {
+                catalog.close();
+            }
+        }
+    }
+
+    private void appendRecord(Table table, int id, OffsetDateTime timestamp) 
throws Exception {
+        appendRecords(table, id, new OffsetDateTime[] {timestamp});
+    }
+
+    private void appendRecords(Table table, int id, OffsetDateTime[] 
timestamps) throws Exception {
+        DataWriter<Record> writer =
+                new GenericAppenderFactory(table.schema(), table.spec())
+                        .newDataWriter(
+                                EncryptedFiles.plainAsEncryptedOutput(
+                                        table.io()
+                                                .newOutputFile(
+                                                        table.location()
+                                                                + "/data/"
+                                                                + id
+                                                                + ".parquet")),
+                                FileFormat.PARQUET,
+                                null);
+        try (DataWriter<Record> closeable = writer) {
+            for (OffsetDateTime timestamp : timestamps) {
+                GenericRecord record = GenericRecord.create(table.schema());
+                record.setField("id", id++);
+                record.setField("event_time", timestamp);
+                closeable.write(record);
+            }
+        }
+        table.newAppend().appendFile(writer.toDataFile()).commit();
+    }
+
+    private Set<Integer> readIds(Table table) throws Exception {
+        Set<Integer> ids = new HashSet<>();
+        try (CloseableIterable<Record> records = 
IcebergGenerics.read(table).build()) {
+            for (Record record : records) {
+                ids.add((Integer) record.getField("id"));
+            }
+        }
+        return ids;
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-iceberg/src/test/java/org/apache/seatunnel/connectors/seatunnel/iceberg/utils/ExpressionUtilsTest.java
 
b/seatunnel-connectors-v2/connector-iceberg/src/test/java/org/apache/seatunnel/connectors/seatunnel/iceberg/utils/ExpressionUtilsTest.java
index 39f3728b94..75cedc2ece 100644
--- 
a/seatunnel-connectors-v2/connector-iceberg/src/test/java/org/apache/seatunnel/connectors/seatunnel/iceberg/utils/ExpressionUtilsTest.java
+++ 
b/seatunnel-connectors-v2/connector-iceberg/src/test/java/org/apache/seatunnel/connectors/seatunnel/iceberg/utils/ExpressionUtilsTest.java
@@ -18,9 +18,12 @@
 package org.apache.seatunnel.connectors.seatunnel.iceberg.utils;
 
 import org.apache.iceberg.Schema;
+import org.apache.iceberg.data.GenericRecord;
+import org.apache.iceberg.expressions.Evaluator;
 import org.apache.iceberg.expressions.Expression;
 import org.apache.iceberg.expressions.Expressions;
 import org.apache.iceberg.types.Types;
+import org.apache.iceberg.util.DateTimeUtil;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
@@ -30,12 +33,111 @@ import net.sf.jsqlparser.parser.CCJSqlParserUtil;
 import net.sf.jsqlparser.statement.Statement;
 import net.sf.jsqlparser.statement.delete.Delete;
 
+import java.time.OffsetDateTime;
+import java.time.format.DateTimeParseException;
 import java.util.List;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 
 public class ExpressionUtilsTest {
 
+    @Test
+    void testTimestampComparisonBoundariesBeforeEpoch() throws Exception {
+        Schema schema =
+                new Schema(
+                        Types.NestedField.optional(
+                                1, "event_time", 
Types.TimestampType.withZone()));
+        String[] operators = {"=", "!=", ">", ">=", "<", "<="};
+        boolean[][] expected = {
+            {false, true, false},
+            {true, false, true},
+            {false, false, true},
+            {false, true, true},
+            {true, false, false},
+            {true, true, false}
+        };
+        for (int i = 0; i < operators.length; i++) {
+            Expression expression =
+                    timestampExpression(
+                            "event_time " + operators[i] + " '1970-01-01 
05:29:59.999999+05:30'",
+                            true);
+            Evaluator evaluator = new Evaluator(schema.asStruct(), expression);
+            for (int j = 0; j < 3; j++) {
+                GenericRecord row = GenericRecord.create(schema);
+                row.setField("event_time", (long) j - 2);
+                assertEquals(
+                        expected[i][j],
+                        evaluator.eval(row),
+                        operators[i] + " at " + ((long) j - 2));
+            }
+        }
+    }
+
+    @Test
+    void testTimestampWithZoneLiterals() throws Exception {
+        long expected =
+                DateTimeUtil.microsFromTimestamptz(
+                        OffsetDateTime.parse("2026-09-12T04:30:00.123456Z"));
+        for (String literal :
+                new String[] {
+                    "2026-09-12 04:30:00.123456Z",
+                    "2026-09-12 10:00:00.123456+05:30",
+                    "2026-09-12 10:00:00.123456999+05:30",
+                    "2026-09-11 23:30:00.123456-05:00",
+                    "2026-09-12 04:30:00.123456"
+                }) {
+            assertEquals(
+                    Expressions.equal("event_time", expected).toString(),
+                    timestampExpression("event_time = '" + literal + "'", 
true).toString());
+        }
+    }
+
+    @Test
+    void testTimestampPredicatesRetainSchema() throws Exception {
+        long expected =
+                DateTimeUtil.microsFromTimestamptz(
+                        OffsetDateTime.parse("2026-09-12T04:30:00.123456Z"));
+        String literal = "'2026-09-12 10:00:00.123456+05:30'";
+        assertEquals(
+                Expressions.in("event_time", expected).toString(),
+                timestampExpression("event_time IN (" + literal + ")", 
true).toString());
+        assertEquals(
+                Expressions.notIn("event_time", expected).toString(),
+                timestampExpression("event_time NOT IN (" + literal + ")", 
true).toString());
+    }
+
+    @Test
+    void testTimestampWithoutZoneRejectsOffsets() {
+        Assertions.assertThrows(
+                DateTimeParseException.class,
+                () -> timestampExpression("event_time = '2026-09-12 
10:00:00+05:30'", false));
+    }
+
+    @Test
+    void testInvalidTimestampWithZoneRejected() {
+        for (String literal :
+                new String[] {
+                    "invalid", "2026-09-12 10:00:00+25:00", "2026-09-12 
10:00:00+05:30 trailing"
+                }) {
+            Assertions.assertThrows(
+                    DateTimeParseException.class,
+                    () -> timestampExpression("event_time = '" + literal + 
"'", true));
+        }
+    }
+
+    private Expression timestampExpression(String predicate, boolean withZone) 
throws Exception {
+        Delete delete = (Delete) CCJSqlParserUtil.parse("DELETE FROM events 
WHERE " + predicate);
+        Schema schema =
+                new Schema(
+                        Types.NestedField.optional(
+                                1,
+                                "event_time",
+                                withZone
+                                        ? Types.TimestampType.withZone()
+                                        : Types.TimestampType.withoutZone()));
+        return ExpressionUtils.convert(delete.getWhere(), schema);
+    }
+
     @Test
     public void testSqlToExpression() throws JSQLParserException {
         String sql = "delete from test.a where id = 1";

Reply via email to