This is an automated email from the ASF dual-hosted git repository.

davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new 244e580440 [Fix][Connector-V2] Support text format in Pulsar source 
(#11792)
244e580440 is described below

commit 244e580440a0dd4c0e61dece1dfe393ffb66c718
Author: 1322630531 <[email protected]>
AuthorDate: Thu Aug 20 23:13:32 2026 +0800

    [Fix][Connector-V2] Support text format in Pulsar source (#11792)
    
    Co-authored-by: DanielLeens <[email protected]>
    Co-authored-by: davidzollo <[email protected]>
---
 docs/en/connectors/source/Pulsar.md                | 29 +++++++++-
 docs/zh/connectors/source/Pulsar.md                | 29 +++++++++-
 .../pulsar/config/PulsarMultiTableConfig.java      | 12 +++-
 .../seatunnel/pulsar/source/PulsarSource.java      | 15 ++++-
 .../pulsar/config/PulsarMultiTableConfigTest.java  | 36 ++++++++++++
 .../seatunnel/pulsar/source/PulsarSourceTest.java  | 65 ++++++++++++++++++++++
 6 files changed, 178 insertions(+), 8 deletions(-)

diff --git a/docs/en/connectors/source/Pulsar.md 
b/docs/en/connectors/source/Pulsar.md
index be83ebcbe8..72231eedf3 100644
--- a/docs/en/connectors/source/Pulsar.md
+++ b/docs/en/connectors/source/Pulsar.md
@@ -41,7 +41,7 @@ Source connector for Apache Pulsar.
 | cursor.stop.mode         | Enum    | No       | NEVER         | Stop 
position mode. Options: `NEVER` (streaming), `LATEST` (batch), `TIMESTAMP` 
(batch)                         |
 | cursor.stop.timestamp    | Long    | No       | -             | Stop 
timestamp (ms) when `cursor.stop.mode=TIMESTAMP`                                
                            |
 | schema                   | Config  | No       | -             | Data 
structure including field names and types                                       
                            |
-| format                   | String  | No       | json          | Data format. 
Default is json. Supported formats: json, canal_json and avro. **Multi-table 
mode only supports JSON, CANAL_JSON and AVRO**                            |
+| format                   | String  | No       | json          | Data format. 
Default is json. Supported formats: json, canal_json, avro and text. **Text is 
supported only in single-table mode; multi-table mode supports JSON, CANAL_JSON 
and AVRO** |
 | field_delimiter          | String  | No       | ,             | Field 
delimiter for `text` format.                                                    
                           |
 | common-options           |         | No       | -             | Source 
plugin common parameters. See [Source Common 
Options](../common-options/source-common-options.md) for details           |
 
@@ -163,7 +163,7 @@ reference to 
[Schema-Feature](../../introduction/concepts/schema-feature.md)
 
 ### format [String]
 
-Data format. The default format is json. Supported formats are json, 
canal_json and avro. The `schema` option is required when using avro format. 
See [formats](../formats) for more details.
+Data format. The default format is json. Supported formats are json, 
canal_json, avro and text. The `schema` option is required when using avro 
format. Text format is supported only in single-table mode. See 
[formats](../formats) for more details.
 
 ### field_delimiter [String]
 
@@ -206,6 +206,31 @@ source {
 }
 ```
 
+### Read Text Messages
+
+Use `format = text` in single-table mode to split each message into schema 
fields with `field_delimiter`.
+
+```hocon
+source {
+  Pulsar {
+    topic = "text-events"
+    subscription.name = "seatunnel-text-sub"
+    client.service-url = "pulsar://localhost:6650"
+    admin.service-url = "http://localhost:8080";
+    cursor.startup.mode = "EARLIEST"
+    cursor.stop.mode = "LATEST"
+    format = text
+    field_delimiter = "|"
+    schema = {
+      fields {
+        id = int
+        name = string
+      }
+    }
+  }
+}
+```
+
 ### Read Canal JSON Messages
 
 Use `format = canal_json` when the Pulsar topic stores Canal JSON change 
events.
diff --git a/docs/zh/connectors/source/Pulsar.md 
b/docs/zh/connectors/source/Pulsar.md
index 55700815c8..bd31b64437 100644
--- a/docs/zh/connectors/source/Pulsar.md
+++ b/docs/zh/connectors/source/Pulsar.md
@@ -41,7 +41,7 @@ Apache Pulsar 的源连接器。
 | cursor.stop.mode         | Enum    | 否    | NEVER  | 
停止位置模式。可选值:`NEVER`(流式)、`LATEST`(批式)、`TIMESTAMP`(批式)                             
          |
 | cursor.stop.timestamp    | Long    | 否    | -      | 当 
`cursor.stop.mode=TIMESTAMP` 时的停止时间戳(毫秒)                                        
        |
 | schema                   | Config  | 否    | -      | 数据结构,包括字段名称和字段类型        
                                                                  |
-| format                   | String  | 否    | json   | 数据格式。默认为 json。支持 
json、canal_json 和 avro 格式。**多表模式仅支持 JSON、CANAL_JSON 和 AVRO**                    
                           |
+| format                   | String  | 否    | json   | 数据格式。默认为 json。支持 
json、canal_json、avro 和 text 格式。**text 仅支持单表模式;多表模式支持 JSON、CANAL_JSON 和 AVRO**   
                   |
 | field_delimiter          | String  | 否    | ,      | `text` 格式使用的字段分隔符。      
                                                                       |
 | common-options           |         | 否    | -      | Source 插件通用参数,请参考 
[Source Common Options](../common-options/source-common-options.md) 了解详情        
       |
 
@@ -156,7 +156,7 @@ Pulsar 消费者的启动模式,有效值为 `'EARLIEST'`、`'LATEST'`、`'SUB
 
 ### format [String]
 
-数据格式。默认值为 `json`。支持 json、canal_json 和 avro 格式。使用 avro 格式时需要配置 
`schema`。更多格式说明参考 [formats](../formats)。
+数据格式。默认值为 `json`。支持 json、canal_json、avro 和 text 格式。使用 avro 格式时需要配置 
`schema`。text 格式仅支持单表模式。更多格式说明参考 [formats](../formats)。
 
 ### field_delimiter [String]
 
@@ -199,6 +199,31 @@ source {
 }
 ```
 
+### 读取 Text 消息
+
+在单表模式下使用 `format = text`,通过 `field_delimiter` 将每条消息拆分为 schema 中定义的字段。
+
+```hocon
+source {
+  Pulsar {
+    topic = "text-events"
+    subscription.name = "seatunnel-text-sub"
+    client.service-url = "pulsar://localhost:6650"
+    admin.service-url = "http://localhost:8080";
+    cursor.startup.mode = "EARLIEST"
+    cursor.stop.mode = "LATEST"
+    format = text
+    field_delimiter = "|"
+    schema = {
+      fields {
+        id = int
+        name = string
+      }
+    }
+  }
+}
+```
+
 ### 读取 Canal JSON 消息
 
 当 Pulsar topic 中保存的是 Canal JSON 变更事件时,使用 `format = canal_json`。
diff --git 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarMultiTableConfig.java
 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarMultiTableConfig.java
index e17db3ddcd..dd4bd7141c 100644
--- 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarMultiTableConfig.java
+++ 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarMultiTableConfig.java
@@ -320,13 +320,21 @@ public class PulsarMultiTableConfig implements 
Serializable {
                         ? String.format("tables_configs[%d] ('%s')", index, 
tablePath)
                         : "Pulsar source config";
         String normalized = format.toUpperCase();
+        if (index >= 0 && Objects.equals("TEXT", normalized)) {
+            throw new PulsarConnectorException(
+                    SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED,
+                    String.format(
+                            "%s uses format 'TEXT', but TEXT is only supported 
in single-table mode",
+                            configPrefix));
+        }
         if (!Objects.equals("JSON", normalized)
                 && !Objects.equals("CANAL_JSON", normalized)
-                && !Objects.equals("AVRO", normalized)) {
+                && !Objects.equals("AVRO", normalized)
+                && !Objects.equals("TEXT", normalized)) {
             throw new PulsarConnectorException(
                     SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED,
                     String.format(
-                            "%s uses unsupported format '%s', only JSON, 
CANAL_JSON and AVRO are supported",
+                            "%s uses unsupported format '%s', only JSON, 
CANAL_JSON, AVRO and TEXT are supported",
                             configPrefix, format));
         }
     }
diff --git 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSource.java
 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSource.java
index f4c28ff6fe..5de647d83d 100644
--- 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSource.java
+++ 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSource.java
@@ -56,6 +56,7 @@ import 
org.apache.seatunnel.format.avro.AvroDeserializationSchema;
 import org.apache.seatunnel.format.json.JsonDeserializationSchema;
 import org.apache.seatunnel.format.json.canal.CanalJsonDeserializationSchema;
 import org.apache.seatunnel.format.json.exception.SeaTunnelJsonFormatException;
+import org.apache.seatunnel.format.text.TextDeserializationSchema;
 
 import org.apache.pulsar.shade.org.apache.commons.lang3.StringUtils;
 
@@ -199,7 +200,7 @@ public class PulsarSource
                     new PulsarConsumerMetadata(
                             tablePath,
                             catalogTable,
-                            createDeserialization(tableConfig.getFormat(), 
catalogTable),
+                            createDeserialization(tableConfig, catalogTable),
                             createDiscoverer(tableConfig),
                             createStartCursor(tableConfig),
                             createStopCursor(tableConfig),
@@ -262,7 +263,8 @@ public class PulsarSource
     }
 
     private DeserializationSchema<SeaTunnelRow> createDeserialization(
-            String format, CatalogTable catalogTable) {
+            PulsarTableConfig tableConfig, CatalogTable catalogTable) {
+        String format = tableConfig.getFormat();
         switch (format.toUpperCase()) {
             case "JSON":
                 return new JsonDeserializationSchema(
@@ -274,6 +276,15 @@ public class PulsarSource
                                 .build());
             case "AVRO":
                 return new AvroDeserializationSchema(catalogTable);
+            case "TEXT":
+                return TextDeserializationSchema.builder()
+                        .seaTunnelRowType(catalogTable.getSeaTunnelRowType())
+                        .delimiter(
+                                tableConfig
+                                        .getSchemaConfig()
+                                        
.get(PulsarSourceOptions.FIELD_DELIMITER))
+                        .setCatalogTable(catalogTable)
+                        .build();
             default:
                 throw new SeaTunnelJsonFormatException(
                         CommonErrorCode.UNSUPPORTED_DATA_TYPE, "Unsupported 
format: " + format);
diff --git 
a/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarMultiTableConfigTest.java
 
b/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarMultiTableConfigTest.java
index 0201493603..7062f5cc97 100644
--- 
a/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarMultiTableConfigTest.java
+++ 
b/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarMultiTableConfigTest.java
@@ -73,6 +73,42 @@ class PulsarMultiTableConfigTest {
         Assertions.assertEquals("AVRO", 
multiTableConfig.getTableConfigs().get(0).getFormat());
     }
 
+    @Test
+    void shouldAllowTextFormatInSingleTable() {
+        Map<String, Object> config = createBaseConfig();
+        config.put("topic", "persistent://public/default/events");
+        config.put("format", "text");
+        config.put("field_delimiter", "|");
+
+        PulsarMultiTableConfig multiTableConfig =
+                PulsarMultiTableConfig.of(ReadonlyConfig.fromMap(config));
+
+        PulsarTableConfig tableConfig = 
multiTableConfig.getTableConfigs().get(0);
+        Assertions.assertEquals("text", tableConfig.getFormat());
+        Assertions.assertEquals(
+                "|", 
tableConfig.getSchemaConfig().get(PulsarSourceOptions.FIELD_DELIMITER));
+    }
+
+    @Test
+    void shouldRejectTextFormatInTablesConfigs() {
+        Map<String, Object> config = createBaseConfig();
+        config.put(
+                "tables_configs",
+                Collections.singletonList(
+                        createTableConfig(
+                                "db.events", 
"persistent://public/default/events", null, "text")));
+
+        PulsarConnectorException exception =
+                Assertions.assertThrows(
+                        PulsarConnectorException.class,
+                        () -> 
PulsarMultiTableConfig.of(ReadonlyConfig.fromMap(config)));
+
+        Assertions.assertEquals(
+                SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED, 
exception.getSeaTunnelErrorCode());
+        Assertions.assertTrue(
+                exception.getMessage().contains("TEXT is only supported in 
single-table mode"));
+    }
+
     @Test
     void shouldRejectOverlappingTopicDeclarations() {
         Map<String, Object> config = createBaseConfig();
diff --git 
a/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSourceTest.java
 
b/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSourceTest.java
index cc3db70982..c3d5b6fcaf 100644
--- 
a/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSourceTest.java
+++ 
b/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSourceTest.java
@@ -19,22 +19,55 @@ package 
org.apache.seatunnel.connectors.seatunnel.pulsar.source;
 
 import org.apache.seatunnel.api.common.JobContext;
 import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.serialization.DeserializationSchema;
 import org.apache.seatunnel.api.source.Boundedness;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
 import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.api.table.catalog.Column;
+import org.apache.seatunnel.api.table.catalog.PhysicalColumn;
+import org.apache.seatunnel.api.table.catalog.TableIdentifier;
+import org.apache.seatunnel.api.table.catalog.TablePath;
+import org.apache.seatunnel.api.table.catalog.TableSchema;
 import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
 import org.apache.seatunnel.common.constants.JobMode;
 import org.apache.seatunnel.common.utils.SerializationUtils;
 import 
org.apache.seatunnel.connectors.seatunnel.pulsar.exception.PulsarConnectorException;
+import org.apache.seatunnel.format.text.TextDeserializationSchema;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
+import java.lang.reflect.Field;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.HashMap;
+import java.util.List;
 import java.util.Map;
 
 class PulsarSourceTest {
 
+    @Test
+    void shouldDeserializeTextWithConfiguredFieldDelimiter() throws Exception {
+        Map<String, Object> config = createBaseConfig("NEVER");
+        config.put("topic", "persistent://public/default/events");
+        config.put("format", "text");
+        config.put("field_delimiter", "|");
+
+        PulsarSource source =
+                new PulsarSource(ReadonlyConfig.fromMap(config), 
createTextCatalogTable());
+        DeserializationSchema<SeaTunnelRow> deserializationSchema =
+                getDeserializationSchema(source);
+
+        Assertions.assertTrue(deserializationSchema instanceof 
TextDeserializationSchema);
+        SeaTunnelRow row =
+                
deserializationSchema.deserialize("1|Alice".getBytes(StandardCharsets.UTF_8));
+        Assertions.assertEquals(1, row.getField(0));
+        Assertions.assertEquals("Alice", row.getField(1));
+    }
+
     @Test
     void shouldExposeProducedCatalogTablesForTablesConfigs() {
         Map<String, Object> config = createBaseConfig("NEVER");
@@ -202,6 +235,38 @@ class PulsarSourceTest {
         return config;
     }
 
+    private CatalogTable createTextCatalogTable() {
+        List<Column> columns = new ArrayList<>();
+        columns.add(
+                PhysicalColumn.builder()
+                        .name("id")
+                        .dataType(BasicType.INT_TYPE)
+                        .nullable(true)
+                        .build());
+        columns.add(
+                PhysicalColumn.builder()
+                        .name("name")
+                        .dataType(BasicType.STRING_TYPE)
+                        .nullable(true)
+                        .build());
+        return CatalogTable.of(
+                TableIdentifier.of("default", "default", "pulsar_table"),
+                TableSchema.builder().columns(columns).build(),
+                new HashMap<>(),
+                new ArrayList<>(),
+                "Pulsar text table");
+    }
+
+    @SuppressWarnings("unchecked")
+    private DeserializationSchema<SeaTunnelRow> 
getDeserializationSchema(PulsarSource source)
+            throws ReflectiveOperationException {
+        Field field = 
PulsarSource.class.getDeclaredField("consumerMetadataMap");
+        field.setAccessible(true);
+        Map<TablePath, PulsarConsumerMetadata> consumerMetadataMap =
+                (Map<TablePath, PulsarConsumerMetadata>) field.get(source);
+        return 
consumerMetadataMap.values().iterator().next().getDeserializationSchema();
+    }
+
     private Map<String, Object> createTableConfig(String tablePath, String 
topic) {
         Map<String, Object> config = new HashMap<>();
         config.put("table_path", tablePath);

Reply via email to