TimurRakhmatullin86 opened a new pull request, #40275:
URL: https://github.com/apache/beam/pull/40275

   ## Summary
   
   When a thread catches `InterruptedException`, the JVM clears the thread's 
interrupt flag. Per Java concurrency best practices ([Java Concurrency in 
Practice, Ch. 7](https://jcip.net/)), the catch block must either re-throw the 
`InterruptedException` or restore the flag via 
`Thread.currentThread().interrupt()`. Failing to do so silently swallows the 
interrupt signal, which can cause the thread (and its callers) to miss shutdown 
or cancellation requests — a particularly dangerous bug in a distributed data 
processing framework like Beam where pipeline cancellation relies on interrupt 
propagation.
   
   This PR fixes **33 occurrences across 27 files** where 
`InterruptedException` was caught without restoring the interrupt status.
   
   ## Affected modules
   
   | Area | Files |
   |------|-------|
   | **Core SDK** | `TransformUpgrader`, `BeamFnDataGrpcMultiplexer` |
   | **I/O Connectors** | Kafka (`KafkaUnboundedReader`), Kinesis 
(`ShardReadersPool`), BigQuery (`BigQueryStorageStreamSource`, 
`DynamicDestinationsHelpers`), Spanner (`SpannerIO`), FHIR (`FhirIO`), Hadoop 
(`HadoopFormatIO`), Spark Receiver (`ReadFromSparkReceiverWithOffsetDoFn`), 
Request-Response (`Call`, `Repeater`) |
   | **Runners** | Flink (`SplittableDoFnOperator`, `StreamingImpulseSource` 
x2), Jet (`JetPipelineResult`), Prism (`PrismExecutor`), fn-execution 
(`GrpcDataService`, `DockerCommand`, `BeamWorkerStatusGrpcService`) |
   | **Extensions** | Python (`PythonService`, `PythonExternalTransform`), GCP 
(`RetryHttpRequestInitializer`) |
   | **Harness** | `FinalizeBundleHandler` |
   | **Transform Service** | `ExpansionService`, `TransformServiceLauncher` |
   | **Testing** | Nexmark (`NexmarkLauncher` — used `Thread.interrupted()` 
which *clears* the flag instead of `Thread.currentThread().interrupt()` which 
*sets* it) |
   
   ## Categories of fixes
   
   1. **Silent swallow** — catch blocks with `// ignore` comment that discarded 
the interrupt entirely
   2. **Multi-catch conflation** — `catch (IOException | InterruptedException 
e)` blocks that handled `InterruptedException` identically to `IOException`, 
losing the interrupt semantics. Fixed by adding `if (e instanceof 
InterruptedException) Thread.currentThread().interrupt();`
   3. **Wrap-and-throw without restore** — catch blocks that wrapped in 
`RuntimeException`/`IllegalStateException` and threw, but didn't restore the 
interrupt flag first
   4. **Wrong API** — `NexmarkLauncher` called `Thread.interrupted()` (a static 
method that reads and *clears* the flag) instead of 
`Thread.currentThread().interrupt()` (which *sets* the flag)
   
   ## Test plan
   
   - [ ] Existing unit tests continue to pass (these changes only add 
`Thread.currentThread().interrupt()` calls, which is a no-op if the thread is 
not being interrupted)
   - [ ] No behavioral change for non-interrupted threads
   - [ ] Interrupted threads now correctly propagate their interrupt status to 
callers
   
   🤖 Generated with [Claude Code](https://claude.com/claude-code)


-- 
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]

Reply via email to