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]

Reply via email to