open-telemetry / open-telemetry/opentelemetry-python-contrib
ProxiedConsumer is creating spans each time it polls instead of only when it receives a message
Nobody has claimed this yet.
- Dominant language
- Python
- Stars
- 1.1k
- Forks
- 1.1k
- Avg merge
- 4d 15h
- Merged PRs (30d)
- 16
Description
Describe your environment
Python 3.8
confluent-kafka 1.8.2
opentelemetry-instrumentation-confluent-kafka 0.37b0
Steps to reproduce
Called the poll method on the ProxiedConsumer from a while true loop, as it's the recommended usage from the confluent-kafka Consumer example given here: https://github.com/confluentinc/confluent-kafka-python#basic-consumer-example
What is the expected behavior?
The poll method on ProxiedConsumer calls wrap_poll method, which should:
- First ensure that the
recordreturned from calling the underlying confluent-kafka.Consumer's poll method exists. - After checking that the record exists, it should extract the context from the record headers, and then start a span using this context, so that the span stays linked to the spans that have been created before the consumer received the current kafka message.
What is the actual behavior?
- The
wrap_pollmethod is creating a span each time it is called. It starts this span before extracting context from the kafka message, so this span is no longer linked to any previous spans. Whereas it should only create a span after checking that the received record is not None (here) and is an actual kafka message. - Since the span started before checking if record exists here is started as current span, the span that's started after record is received will use the current span's context even if the links contain the context from the message headers.
- Lastly,
wrap_pollreturns the record even if the record is None, it should only return record if the record exists.
Additional context
I want to confirm if the way I'm using ProxiedConsumer and its poll method is correct. Here's my understanding and how I'm using it:
ProxiedConsumer is a wrapper around the confluent-kafka Consumer class. ProxiedConsumer has a method call poll, which calls ConfluentKafkaInstrumentor's wrap_poll method. And wrap_poll method calls the underlying Consumer's poll method with the user specified timeout.
As a user of opentelemetry-instrumentation-confluent-kafka, one would create the ProxiedConsumer and use it as follows based on the docs:
c = confluent_kafka.Consumer({ 'bootstrap.servers': 'localhost:29092' })
consumer = ConfluentKafkaInstrumentor().instrument_consumer(c, tracer_provider=tracer_provider)
consumer.subscribe(['mytopic'])
while True:
msg = consumer.poll(1.0)
if msg is None:
continue
process(msg) // process is just some method to handle messages
As in the code above, the ProxiedConsumer's poll would be called from a while True loop, with a timeout of 1, resulting in spans per second, even if the msg from kafka could be None.
- The
wrap_pollmethod should instead start the span only if the record it gets from callingpollactually exists. - It should start this span after extracting context from the record headers and using that as the context while starting span as current span.
- After it starts the span, it should also inject the current context in the record headers, so that future spans for this record are created using the current span's tracing headers
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 instrumentation/opentelemetry-instrumentation-confluent-kafka/src/opentelemetry/instrumentation/confluent_kafka/init.py at wrap_poll and inspect the underlying Consumer.poll call. Confirm behavior with the ProxiedConsumer usage described in the issue: no span for a None record, context extraction before starting a span, and returning only an existing record.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- kafka, python
- Domain
- distributed-systems
- Issue type
- Bug
- Difficulty
- 3/5
- Estimated time
- 1-2 days
- Activity status
- Quiet
- Clarity
- Clearly specified
- Newbie friendliness
- 72/100