Skip to content

SynchronousMultiSpanProcessor.force_flush stops when a processor doesn't return True #4631

Description

@alexmojaki

This code in SynchronousMultiSpanProcessor:

def force_flush(self, timeout_millis: int = 30000) -> bool:
"""Sequentially calls force_flush on all underlying
:class:`SpanProcessor`
Args:
timeout_millis: The maximum amount of time over all span processors
to wait for spans to be exported. In case the first n span
processors exceeded the timeout followup span processors will be
skipped.
Returns:
True if all span processors flushed their spans within the
given timeout, False otherwise.
"""
deadline_ns = time_ns() + timeout_millis * 1000000
for sp in self._span_processors:
current_time_ns = time_ns()
if current_time_ns >= deadline_ns:
return False
if not sp.force_flush((deadline_ns - current_time_ns) // 1000000):
return False

means that if a processor doesn't return something truthy, then remaining processors are skipped. This doesn't really match the docstring, which only mentions skipping in case of exceeding the timeout. To be fair, SpanProcessor.force_flush does say Returns: False if the timeout is exceeded, True otherwise. But it would be safer to just check the timeout in SynchronousMultiSpanProcessor. Why skip the rest otherwise?

This is particularly bad when one of the processors doesn't implement force_flush, so it inherits the default implementation which just returns None, which is interpreted as exceeding the timeout. To demonstrate:

from opentelemetry.sdk.trace import SpanProcessor, TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter


class MySpanProcessor(SpanProcessor):
    def on_start(self, span, parent_context):
        print(f'Span started: {span.name}')

    # force_flush not implemented


exporter1 = InMemorySpanExporter()
exporter2 = InMemorySpanExporter()
batch_processor_1 = BatchSpanProcessor(exporter1)
batch_processor_2 = BatchSpanProcessor(exporter2)

tracer_provider = TracerProvider()
tracer_provider.add_span_processor(batch_processor_1)
tracer_provider.add_span_processor(MySpanProcessor())  # has default empty force_flush which returns None
tracer_provider.add_span_processor(batch_processor_2)  # this will get skipped because of the previous processor

tracer = tracer_provider.get_tracer('my_tracer')

tracer.start_span('my_span').end()

tracer_provider.force_flush()

assert len(exporter1.get_finished_spans()) == 1
assert len(exporter2.get_finished_spans()) == 0  # skipped!!!

If for some reason it's important to skip other processors after one returns False, then I still think that:

  1. SpanProcessor.force_flush should return True by default
  2. SynchronousMultiSpanProcessor should distinguish between False and None.

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Fields

    No fields configured for issues without a type.

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions