Dream95 opened a new issue, #11620: URL: https://github.com/apache/seatunnel/issues/11620
### Search before asking - [x] I had searched in the [issues](https://github.com/apache/seatunnel/issues?q=is%3Aissue+label%3A%22bug%22) and found no similar issues. ### What happened Pulsar Source option contract and factory OptionRule claim that `format = text` is supported, but the runtime validation and deserialization path reject it. **Contract / docs mismatch:** 1. `PulsarBaseOptions.FORMAT` description says: `Optional text and avro format` 2. `PulsarSourceFactory` OptionRule binds `field_delimiter` when `format = text` 3. Source docs still document `field_delimiter` for text format **Runtime behavior:** 1. `PulsarMultiTableConfig.validateFormat()` only allows `JSON`, `CANAL_JSON` and `AVRO` (also applied to single-table configs) 2. `PulsarSource.createDeserialization()` has no `TEXT` branch and throws `Unsupported format: text` As a result, configuring Pulsar Source with `format = text` fails during config validation / source initialization. Note: Pulsar **Sink** already supports text format via `TextSerializationSchema`. The dependency `seatunnel-format-text` is already present in `connector-pulsar`. ### SeaTunnel Version dev (master) ### SeaTunnel Config ```conf env { parallelism = 1 job.mode = "BATCH" } source { Pulsar { topic = "persistent://public/default/test-topic" client.service-url = "pulsar://localhost:6650" admin.service-url = "http://localhost:8080" subscription.name = "seatunnel-sub" cursor.startup.mode = "EARLIEST" cursor.stop.mode = "LATEST" format = text field_delimiter = "," schema = { fields { id = "int" name = "string" } } } } sink { Console {} } ``` ### Running Command ```shell ./bin/seatunnel.sh --config config/pulsar_text_source.conf -e local ``` ### Error Exception ```log org.apache.seatunnel.connectors.seatunnel.pulsar.exception.PulsarConnectorException: Pulsar source config uses unsupported format 'text', only JSON, CANAL_JSON and AVRO are supported ``` If validation were bypassed, deserialization would still fail with: ```log Unsupported format: text ``` ### Zeta or Flink or Spark Version Zeta (local) ### Java or Scala Version Java 8 / 11 ### Screenshots N/A ### Root Cause | Layer | Behavior for `text` | | --- | --- | | `PulsarBaseOptions.FORMAT` | Declares text as optional | | `PulsarSourceFactory` OptionRule | Conditional `field_delimiter` for text | | `PulsarMultiTableConfig.validateFormat()` | Rejects text | | `PulsarSource.createDeserialization()` | No TEXT case | | Pulsar Sink | Text is supported | ### Expected Behavior Either: 1. **Implement text format for Pulsar Source** (recommended, aligns with OptionRule / Sink / `seatunnel-format-text` dependency), including: - Allow `text` in `validateFormat()` for single-table mode - Add `TextDeserializationSchema` in `createDeserialization()` using `field_delimiter` - Keep multi-table restriction explicit if text remains unsupported there - Update `docs/en` and `docs/zh` Source Pulsar docs accordingly or 2. Remove text from Source Option description / OptionRule / docs if text is intentionally unsupported on Source. ### Are you willing to submit PR? - [x] Yes I am willing to submit a PR! ### Code of Conduct - [x] I agree to follow this project's [Code of Conduct](https://www.apache.org/foundation/policies/conduct) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
