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