apache / apache/hudi

[SUPPORT] - Cannot Ingest Protobuf records using Hudi Streamer

Open
#12,301 11 comments 0 reactions 0 assignees View on GitHub
area:ingest
Dominant language
Java
Stars
6.2k
Forks
2.5k
Avg merge
2d 8h
Merged PRs (30d)
111

Description

I’m trying to ingest from a ProtoKafka source using Hudi Streamer but encountering an issue.

```
Exception in thread "main" org.apache.hudi.utilities.ingestion.HoodieIngestionException: Ingestion service was shut down with exception.
at ...
Error reading source schema from registry. Please check hoodie.streamer.schemaprovider.registry.url is configured correctly. Truncated URL: https://....ons/latest
at org.apache.hudi.utilities.schema.SchemaRegistryProvider.parseSchemaFromRegistry(SchemaRegistryProvider.java:111)
at org.apache.hudi.utilities.schema.SchemaRegistryProvider.getSourceSchema(SchemaRegistryProvider.java:204)
... 10 more
...
Caused by: org.apache.hudi.internal.schema.HoodieSchemaException: Failed to parse schema from registry: syntax = "proto3";
package datagen;
...
Caused by: java.lang.NoSuchMethodException: org.apache.hudi.utilities.schema.converter.ProtoSchemaToAvroSchemaConverter.()
at java.lang.Class.getConstructor0(Class.java:3082)
at java.lang.Class.newInstance(Class.java:412)
... 13 more
```
The stack trace points to a misconfigured schema registry URL. However, the same URL works for Hudi streamer jobs ingesting from AvroKafka sources. When I ping the schema registry URL using curl, it correctly returns the schema.

Additional Context

1. I've verified the Protobuf schema is valid, it is a sample proto schema from Confluent’s Datagen connector.
2. I've confirmed the schema registry URL is configured correctly, it works fine with a similar`AvroKafka` spark job.
3. I added `hoodie.streamer.schemaprovider.proto.class.name` and `hoodie.streamer.source.kafka.proto.value.deserializer.class=org.apache.kafka.common.serialization.ByteArrayDeserializer`. I don't think these are required but their presence/absence did not resolve this error.

Environment Details
Hudi version: v0.15.0
Spark version: 3.1.3
Scala version: 2.12
Google Dataproc version: 2.0.125-debian10

Spark Submit Command and Protobuf Configuration
```
gcloud dataproc jobs submit spark --cluster \
--region us-central1 \
--class org.apache.hudi.utilities.streamer.HoodieStreamer \
--project \
--jars /jars/hudi-gcp-bundle-0.15.0.jar,/jars/spark-avro_2.12-3.1.1.jar,/jars/hudi-utilities-bundle-raw_2.12-0.15.0.jar,/jars/kafka-protobuf-provider-5.5.0.jar \
--schemaprovider-class org.apache.hudi.utilities.schema.SchemaRegistryProvider \
--source-class org.apache.hudi.utilities.sources.ProtoKafkaSource \
--hoodie-conf sasl.jaas.config="org.apache.kafka.common.security.plain.PlainLoginModule required username='' password='';" \
--hoodie-conf hoodie.streamer.schemaprovider.proto.class.name= \
--hoodie-conf basic.auth.credentials.source=USER_INFO \
--hoodie-conf schema.registry.basic.auth.user.info=: \
--hoodie-conf hoodie.streamer.schemaprovider.registry.url=https://:@/subjects/-value/versions/latest \
--hoodie-conf hoodie.streamer.source.kafka.topic= \
--hoodie-conf hoodie.streamer.source.kafka.value.deserializer.class=io.confluent.kafka.serializers.protobuf.KafkaProtobufDeserializer \
--hoodie-conf hoodie.streamer.schemaprovider.registry.schemaconverter=org.apache.hudi.utilities.schema.converter.ProtoSchemaToAvroSchemaConverter \
```

Steps to Reproduce

1. Build a Hudi 0.15.0 JAR with Spark 3.1 and Scala 2.12.
2. Use a Protobuf schema on an accessible schema registry, preferably an authenticated one.
3. Configure Hudi Streamer job with the spark submit command above.
4. Run the Spark job.

I’d appreciate any insights into resolving this issue.
Is there an alternative or a workaround for configuring the Protobuf schema?
Am I missing any configuration settings?
Thank you for your help!

Contributor guide

No contributing guide indexed for this repository

Research direction

Start by reproducing the Spark job with the listed Hudi Streamer configuration, focusing on SchemaRegistryProvider.java:111 and :204 and ProtoSchemaToAvroSchemaConverter. Check how ProtoKafkaSource and the registry schema converter are loaded with the supplied dependencies and configuration. Done means a valid Protobuf schema is read from the registry and records are ingested successfully.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, kafka, spark
Domain
data-engineering, stream-processing
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
28/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.