Skip to content

Span Buffer Multiprocess Enhancement with Health Monitoring - #1830

Open
kaihao-zhao wants to merge 2 commits into
span-flusher-stablefrom
span-flusher-multiprocess-hesol61-r1
Open

kaihao-zhao wants to merge 2 commits into
span-flusher-stablefrom
span-flusher-multiprocess-hesol61-r1

Conversation

@kaihao-zhao

Copy link
Copy Markdown

See title.

@unblocked-local-kaihao unblocked-local-kaihao Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

9 issues found.

About Unblocked

Unblocked has been set up to automatically review your team's pull requests to identify genuine bugs and issues.

📖 Documentation — Learn more in our docs.

💬 Ask questions — Mention @unblocked-local-kaihao to request a review or summary, or ask follow-up questions.

👍 Give feedback — React to comments with 👍 or 👎 to help us improve.

⚙️ Customize — Adjust settings in your preferences.

Comment on lines +254 to +255
if isinstance(process, multiprocessing.Process):
process.kill()

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Workers are constructed with multiprocessing.get_context("spawn").Process, whose SpawnProcess class is not a subclass of multiprocessing.Process. This check is therefore false for every production worker. A hung worker survives while its replacement starts flushing the same shards, and its handle is overwritten.

The same incorrect check appears at line 346, preventing shutdown termination. Use self.mp_context.Process or multiprocessing.process.BaseProcess at both sites.

except (ValueError, AttributeError):
pass # Process already closed, ignore

self._create_process_for_shards(process_index, shards)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

With produce_to_pipe, the worker is a thread. If its callback blocks beyond the health timeout, the preceding cleanup does nothing, but this line starts another thread and overwrites the original handle. The replacement can flush the same segments while the original callback is still running; when the original resumes, both workers produce/delete data and update the same health values.

Handle an unhealthy live thread separately: fail the consumer, or stop and join it using a worker-specific cancellation mechanism before replacement.

# Find which process this shard belongs to and restart that process
for process_index, shards in self.process_to_shards_map.items():
if shard in shards:
self._create_process_for_shards(process_index, shards)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Calling _create_process_for_shard() for a shard whose worker is still alive starts another worker without stopping the existing one. _create_process_for_shards() then overwrites the tracked handle, leaving two workers flushing the same shard group and only one available for cleanup.

The helper promises a restart, so it must stop and reap the group's current worker before creating its replacement.

Comment on lines +340 to +341
if remaining_time <= 0:
break

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If next_step.join() consumes the timeout, this breaks before cleaning up any worker. If an earlier worker consumes it, all subsequent workers are skipped. Setting stopped does not stop workers blocked in Redis or future.result(), so those workers can survive shutdown or rebalancing.

Separate deadline-bounded waiting from unconditional cleanup of every worker; expiration should stop waiting, not skip termination.

click.Option(
["--flusher-processes", "flusher_processes"],
default=1,
type=int,

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

--flusher-processes=-1 is accepted and forwarded unchanged. The flusher computes a negative process count, creates an empty shard map, then raises KeyError while assigning the first shard. Zero is also accepted but silently selects one worker per shard through the or fallback.

Require a positive count at the CLI boundary and validate programmatic constructor inputs as well.

Suggested change
type=int,
type=click.IntRange(min=1),

Comment on lines +300 to +301
for buffer in self.buffers.values():
memory_infos.extend(buffer.get_memory_info())

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

SpansBuffer.get_memory_info() returns memory information for the entire configured Redis cluster, not just its assigned shards. All these buffers use that same cluster, so this loop repeats identical cluster-wide INFO queries once per flusher process on every submission when memory backpressure is enabled.

This adds synchronous Redis work proportional to the worker count without providing additional information. Sample the cluster once rather than once per shard buffer.

self.process.kill()
except ValueError:
pass # Process already closed, ignore
if self.process_restarts[process_index] > MAX_PROCESS_RESTARTS:

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

With MAX_PROCESS_RESTARTS = 10, a counter of 10 still passes this check, is incremented to 11, and starts another replacement. The consumer raises only on the following failure, allowing one restart beyond the stated maximum.

Suggested change
if self.process_restarts[process_index] > MAX_PROCESS_RESTARTS:
if self.process_restarts[process_index] >= MAX_PROCESS_RESTARTS:


step.poll()
# Give flusher threads time to process after drift change
time.sleep(0.1)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Line 15 replaces time.sleep with a no-op, so this new call does not give the flusher threads any time despite the preceding comment. It provides no synchronization before join() sets the stop flag, leaving the test dependent on thread scheduling.

Use a saved reference to the real sleep, as the sibling backpressure test does, or synchronize explicitly with the worker.

"""
A background thread that polls Redis for new segments to flush and to produce to Kafka.
A background multiprocessing manager that polls Redis for new segments to flush and to produce to Kafka.
Creates one process per shard for parallel processing.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The implementation shares shards whenever max_processes is below the shard count, and the CLI defaults to one process for all assigned shards. The new “one process per shard” statement contradicts both that default and the new process-limit test.

Suggested change
Creates one process per shard for parallel processing.
Distributes assigned shards across a bounded number of worker processes.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants