opensearch-project / opensearch-project/data-prepper

[BUG] Handle unsupported PROTOBUF schema type gracefully from Confluent Schema registry in kafka source

Open
#4,648 0 comments 0 reactions 0 assignees View on GitHub

Nobody has claimed this yet.

bug
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

Open the contributing guide

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. 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

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.