xiangfu0 commented on code in PR #19737:
URL: https://github.com/apache/pinot/pull/19737#discussion_r4213712163
##########
pinot-segment-local/src/main/java/org/apache/pinot/segment/local/realtime/writer/StatelessRealtimeSegmentWriter.java:
##########
@@ -308,11 +309,8 @@ public void run() {
public void stopConsumption() {
if (_consumerThread.isAlive()) {
_consumerThread.interrupt();
- try {
- _consumerThread.join();
- } catch (InterruptedException e) {
- _logger.warn("Interrupted while waiting for consumer thread to
finish");
- }
+ // Wait even if interrupted, so that the segment is not destroyed while
the consumer thread is still indexing
+ Uninterruptibles.joinUninterruptibly(_consumerThread);
Review Comment:
**[MINOR] [pre-existing, made load-bearing by this change] The consumer loop
never checks the interrupt flag, so this uninterruptible join relies entirely
on the stream plugin throwing on interrupt.**
`PartitionConsumer.run()` exits only when `_currentOffset` reaches
`_endOffset` or when `fetchMessages`/decode/index throws. Kafka's consumer
throws `InterruptException` from a poll on an interrupted thread, but a plugin
whose `fetchMessages` is a bounded poll that swallows the flag (or retries
internally) will just keep consuming to the end offset after
`stopConsumption()` interrupts it. Before this change an interrupted closer
abandoned the join (the unsafe destroy you fixed); now it blocks with no way to
interrupt it.
That matters for the new `@PreDestroy` path: `shutDown()` interrupts a
running worker, which then calls `writer.close()` -> `stopConsumption()` ->
this join. The javadoc claim that running jobs "fail and clean up their
segment" on shutdown therefore holds only for interrupt-responsive plugins;
otherwise the (now daemon) worker hangs here and the segment dirs are not
cleaned until process exit. `testRunningReingestionIsStoppedOnShutdown` cannot
catch this because its mock blocks in `CountDownLatch.await()`, which does
honor the interrupt.
Cheap fix that makes the uninterruptible join safe across plugins: include
`!Thread.currentThread().isInterrupted()` in the `while` condition (treating
interrupt as consumption failure), so the loop terminates deterministically at
the next batch boundary regardless of the plugin's interrupt behavior.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]