opensearch-project / opensearch-project/data-prepper
[BUG] Handle unsupported PROTOBUF schema type gracefully from Confluent Schema registry in kafka source
Nobody has claimed this yet.
- Dominant language
- Java
- Stars
- 374
- Forks
- 354
- Avg merge
- 3d 18h
- Merged PRs (30d)
- 8
Description
Describe the bug
Confluent Schema Registry support PROTOBUF schema, when consuming data from topic with PROTOBUF schema, NPE is detected at runtime.
2024-06-19T17:38:03,209 [main] INFO org.opensearch.dataprepper.pipeline.server.DataPrepperServer - Data Prepper server running at :4900
2024-06-19T17:38:03,506 [kafka-pipeline-sink-worker-2-thread-1] INFO org.opensearch.dataprepper.plugins.kafka.source.KafkaSource - Starting consumer with the properties : {value.deserializer=class org.apache.kafka.common.serialization.StringDeserializer, auto.register.schemas=false, basic.auth.credentials.source=USER_INFO, group.id=topic_3_local7, reconnect.backoff.ms=10000, max.partition.fetch.bytes=1048576, bootstrap.servers=pkc-rgm37.us-west-2.aws.confluent.cloud:9092, retry.backoff.ms=10000, schema.registry.url=https://psrc-e8157.us-east-2.aws.confluent.cloud, enable.auto.commit=false, sasl.mechanism=PLAIN, sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username="2PYKP5A5VDPZDJZU" password="40NGkZ6JYgehKBfDHa2AoLjmIzWwn19HMXthIr/wYankBxPQhFGMI2oQpJ83a0qc";, fetch.max.wait.ms=500, sasl.client.callback.handler.class=class org.opensearch.dataprepper.plugins.kafka.authenticator.DynamicSaslClientCallbackHandler, session.timeout.ms=45000, client.id=AWSOpenSearchIngestion-1540FD4C-C1E6-4F02-880A-4FE8CA226E1E, key.deserializer=org.apache.kafka.common.serialization.StringDeserializer, max.poll.records=500, auto.commit.interval.ms=5000, heartbeat.interval.ms=5000, security.protocol=SASL_SSL, basic.auth.user.info=XRG6T73EHFA32DJU:bxbOivHjQZfw3k88orcxHmeGQz3Xvv5hkwYeQ6urXIPCgC8lxMtiAH9vNrY2py6p, fetch.min.bytes=1, fetch.max.bytes=52428800, max.poll.interval.ms=300000, auto.offset.reset=earliest}
2024-06-19T17:38:03,508 [kafka-pipeline-sink-worker-2-thread-1] ERROR org.opensearch.dataprepper.plugins.kafka.source.KafkaSource - Failed to setup the Kafka Source Plugin.
java.lang.NullPointerException: Cannot invoke "org.opensearch.dataprepper.plugins.kafka.util.MessageFormat.ordinal()" because "schema" is null
at org.opensearch.dataprepper.plugins.kafka.source.KafkaSource.createKafkaConsumer(KafkaSource.java:166) ~[kafka-plugins-2.9.0-SNAPSHOT.jar:?]
at org.opensearch.dataprepper.plugins.kafka.source.KafkaSource.lambda$start$0(KafkaSource.java:129) ~[kafka-plugins-2.9.0-SNAPSHOT.jar:?]
at java.base/java.util.stream.Streams$RangeIntSpliterator.forEachRemaining(Streams.java:104) ~[?:?]
at java.base/java.util.stream.IntPipeline$Head.forEach(IntPipeline.java:617) ~[?:?]
at org.opensearch.dataprepper.plugins.kafka.source.KafkaSource.lambda$start$1(KafkaSource.java:126) ~[kafka-plugins-2.9.0-SNAPSHOT.jar:?]
at java.base/java.util.ArrayList.forEach(ArrayList.java:1511) ~[?:?]
at org.opensearch.dataprepper.plugins.kafka.source.KafkaSource.start(KafkaSource.java:116) ~[kafka-plugins-2.9.0-SNAPSHOT.jar:?]
at org.opensearch.dataprepper.pipeline.Pipeline.startSourceAndProcessors(Pipeline.java:215) ~[data-prepper-core-2.9.0-SNAPSHOT.jar:?]
at org.opensearch.dataprepper.pipeline.Pipeline.lambda$execute$2(Pipeline.java:260) ~[data-prepper-core-2.9.0-SNAPSHOT.jar:?]
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) [?:?]
at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) [?:?]
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) [?:?]
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) [?:?]
at java.base/java.lang.Thread.run(Thread.java:840) [?:?]
2024-06-19T17:38:03,516 [kafka-pipeline-sink-worker-2-thread-1] ERROR org.opensearch.dataprepper.pipeline.common.PipelineThreadPoolExecutor - Pipeline [kafka-pipeline] process worker encountered a fatal exception, cannot proceed further
java.util.concurrent.ExecutionException: java.lang.RuntimeException: java.lang.NullPointerException: Cannot invoke "org.opensearch.dataprepper.plugins.kafka.util.MessageFormat.ordinal()" because "schema" is null
at java.base/java.util.concurrent.FutureTask.report(FutureTask.java:122) ~[?:?]
at java.base/java.util.concurrent.FutureTask.get(FutureTask.java:191) ~[?:?]
at org.opensearch.dataprepper.pipeline.common.PipelineThreadPoolExecutor.afterExecute(PipelineThreadPoolExecutor.java:70) [data-prepper-core-2.9.0-SNAPSHOT.jar:?]
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1137) [?:?]
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) [?:?]
at java.base/java.lang.Thread.run(Thread.java:840) [?:?]
Caused by: java.lang.RuntimeException: java.lang.NullPointerException: Cannot invoke "org.opensearch.dataprepper.plugins.kafka.util.MessageFormat.ordinal()" because "schema" is null
at org.opensearch.dataprepper.plugins.kafka.source.KafkaSource.lambda$start$1(KafkaSource.java:159) ~[kafka-plugins-2.9.0-SNAPSHOT.jar:?]
at java.base/java.util.ArrayList.forEach(ArrayList.java:1511) ~[?:?]
at org.opensearch.dataprepper.plugins.kafka.source.KafkaSource.start(KafkaSource.java:116) ~[kafka-plugins-2.9.0-SNAPSHOT.jar:?]
at org.opensearch.dataprepper.pipeline.Pipeline.startSourceAndProcessors(Pipeline.java:215) ~[data-prepper-core-2.9.0-SNAPSHOT.jar:?]
at org.opensearch.dataprepper.pipeline.Pipeline.lambda$execute$2(Pipeline.java:260) ~[data-prepper-core-2.9.0-SNAPSHOT.jar:?]
Caused by: java.lang.RuntimeException: java.lang.NullPointerException: Cannot invoke "org.opensearch.dataprepper.plugins.kafka.util.MessageFormat.ordinal()" because "schema" is null
To Reproduce
Run data prepper with the kafka source and confluent schema registry with protobuf topic
kafka-pipeline:
source:
kafka:
bootstrap_servers:
- "pkc-rgm37.us-west-2.aws.confluent.cloud:9092"
authentication:
sasl:
plain:
username: <username>
password: <password>
schema:
type: "confluent"
registry_url: "https://psrc-e8157.us-east-2.aws.confluent.cloud"
api_key: <api_key>
api_secret: <api_secret>
basic_auth_credentials_source: "USER_INFO"
encryption:
type: "ssl"
topics:
- name: <protobuf topic>
group_id: <group id>
workers: 1
client_id: "AWSOpenSearchIngestion-1540FD4C-C1E6-4F02-880A-4FE8CA226E1E"
Contributor guide
First steps
- Read the whole issue, then the project's contributing guide.
- Comment on the issue to say you are picking it up — it saves two people doing the same work.
- Fork the repository and make your change on a branch.
- Open a pull request that references the issue number.
Research direction
Start in KafkaSource.java at createKafkaConsumer, especially the schema handling reported at line 166, and trace how the Confluent schema type is mapped before the consumer is created. Reproduce the Kafka source with a Confluent PROTOBUF topic; done means the unsupported schema is handled gracefully without the reported NullPointerException.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- java, kafka
- Domain
- stream-processing
- Issue type
- Bug
- Difficulty
- 3/5
- Estimated time
- 1-2 days
- Activity status
- Stale
- Clarity
- Mostly clear
- Newbie friendliness
- 50/100