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]
