goutamadwant opened a new issue, #12037: URL: https://github.com/apache/seatunnel/issues/12037
## Search before asking - [x] I searched issues and pull requests in all states and found no similar report or active implementation. ## What happened When Kafka uses `format = protobuf` with `strip_schema_registry_header = true`, a valid Confluent Schema Registry message whose Protobuf data section is empty cannot be deserialized. For example, this valid wire-format message fails: `00 00 00 00 01 00` It contains: - Magic byte `0`. - Schema ID `1`. - The optimized one-byte encoding for message-index path `[0]`. - An empty Protobuf data section. An empty data section is valid for an empty Protobuf message and for a proto3 message whose implicit-presence fields all contain default values. Such fields are omitted from the binary encoding. `SchemaRegistryAwareProtobufDeserializationSchema` currently assumes that a valid Protobuf payload must contain at least two bytes. It rejects the empty data section and falls back to deserializing the complete Schema Registry frame as plain Protobuf. That fallback fails with an invalid-tag error. The production path is: Kafka source with `format = protobuf` and `strip_schema_registry_header = true` → `KafkaSourceConfig` → `SchemaRegistryAwareProtobufDeserializationSchema` Expected behavior: A zero-byte Protobuf data section should be passed to the normal Protobuf parser when, and only when, the preceding Confluent magic byte, schema ID, and complete message-index vector are structurally valid. Plain Protobuf fallback and existing non-empty behavior should remain compatible. References: - Confluent wire format: https://docs.confluent.io/platform/current/schema-registry/fundamentals/serdes-develop/overview.html - Protobuf proto3 defaults: https://protobuf.dev/programming-guides/proto3/ ## Why the proposed fix helps The proposed fix structurally parses the Confluent zigzag-varint message-index vector instead of guessing the payload offset. A zero-byte data section is accepted only after the complete header is validated, so valid default-valued messages work without treating malformed framing as an empty message. The existing plain-Protobuf fallback and bounded compatibility probe for non-empty records remain in place. The parser is allocation-free and bounded by the input length, and the change adds no dependency, public API, configuration option, or default-value change. ## SeaTunnel Version Current `dev` at `5dbfb374f985349aeefde9cf84169aa98b3ac5ca`. ## SeaTunnel Config ```conf source { Kafka { topic = "protobuf-default-values" bootstrap.servers = "localhost:9092" start_mode = "earliest" format = protobuf strip_schema_registry_header = true protobuf_message_name = TestMessage protobuf_schema = """ syntax = "proto3"; message TestMessage { int32 id = 1; string name = 2; } """ schema = { fields { id = int name = string } } } } sink { Console {} } ``` The reproduced record uses schema ID `1`, optimized message-index path `[0]`, and a `TestMessage` with `id = 0` and `name = ""`. Its complete bytes are `00 00 00 00 01 00`. ## Running Command The deterministic format-level regression can be run with: ```shell ./mvnw -pl seatunnel-formats/seatunnel-format-protobuf \ -Dtest=SchemaRegistryAwareProtobufDeserializationSchemaTest#shouldDeserializeSchemaRegistryMessageWithDefaultValuedPayload \ test ``` ## Error Exception ```log com.google.protobuf.InvalidProtocolBufferException: Protocol message contained an invalid tag (zero). ``` ## Zeta or Flink or Spark Version Not engine-specific. The failure occurs in the shared Protobuf format deserializer selected by the Kafka source configuration. ## Java or Scala Version Reproduced independently with: - Oracle Java 8 (`1.8.0_172`) - Eclipse Temurin Java 11 (`11.0.19`) The completed Protobuf format module passes on both runtimes: 19 tests, 0 failures, 0 errors on each. ## Proposed fix contract 1. Parse the Confluent message-index vector structurally using bounded zigzag-varint decoding. 2. Permit a zero-byte Protobuf payload only after a complete valid header. 3. Preserve plain Protobuf fallback. 4. Preserve non-empty messages accepted by the existing compatibility probe. 5. Reject malformed, truncated, negative, or oversized index vectors as empty messages. 6. Cover optimized and explicit `[0]`, empty/default-valued messages, nested indexes, malformed headers, plain messages, and legacy non-empty framing. ## Screenshots Not applicable. ## 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. -- 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]
