apache/iceberg · error · RuntimeException

Interrupted while connecting to Zookeeper

Error message

Interrupted while connecting to Zookeeper

What it means

ZkLockFactory.open() throws RuntimeException("Interrupted while connecting to Zookeeper") if the thread is interrupted while waiting for the Curator client to connect (blockUntilConnected). The interrupt flag is restored before throwing, so callers can observe cancellation. This usually means the enclosing task was cancelled during startup.

Solutions

  1. Fix the underlying connectivity problem that made the connect step slow enough to be interrupted (see the ZK address/port/ensemble).
  2. Increase connectionTimeoutMs so startup completes quickly instead of lingering until cancelled.
  3. If the interrupt is expected (job cancel), treat this as normal cancellation and restore the interrupt flag in your caller.
  4. Check rollout/restart logs to correlate the interruption with job cancellation events.

Example fix

// before
Thread.currentThread().interrupt(); // swallowed in user code, job loop spins
// after
try {
  lockFactory.open();
} catch (RuntimeException e) {
  if (e.getCause() instanceof InterruptedException) {
    Thread.currentThread().interrupt(); // honor cancellation
    return;
  }
  throw e;
}
Defensive patterns

Strategy: try-catch

Validate before calling

if (Thread.currentThread().isInterrupted()) { /* skip ZK connect */ }

Try / catch

try { lockFactory.open(); } catch (RuntimeException e) { if (e.getCause() instanceof InterruptedException) { Thread.currentThread().interrupt(); return; } throw e; }

Prevention

When it happens

Trigger: open() blocks in client.blockUntilConnected(...) waiting for the ZK session and the executing thread is interrupted (Flink task cancel, job restart, executor shutdown).

Common situations: User cancels a job that is stuck connecting to an unreachable ZooKeeper; deployment rollouts cancelling tasks mid-startup; a watchdog interrupting threads that hang on a bad ZK address.

Related errors


AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12). Data as JSON: /api/errors/cbd4ba13111aeacd. Report an issue: GitHub.

Appendix: source

Thrown at flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/maintenance/api/ZkLockFactory.java:137

            .connectionTimeoutMs(connectionTimeoutMs)
            .retryPolicy(createRetryPolicy())
            .build();
    client.start();

    try {
      if (!client.blockUntilConnected(connectionTimeoutMs, TimeUnit.MILLISECONDS)) {
        throw new IllegalStateException("Connection to Zookeeper timed out");
      }

      this.taskSharedCount = new SharedCount(client, getTaskSharePath(), 0);
      this.recoverySharedCount = new SharedCount(client, getRecoverySharedPath(), 0);
      taskSharedCount.start();
      recoverySharedCount.start();
      isOpen = true;
      LOG.info("ZkLockFactory initialized for lockId: {}.", lockId);
    } catch (InterruptedException e) {
      Thread.currentThread().interrupt();
      throw new RuntimeException("Interrupted while connecting to Zookeeper", e);
    } catch (Exception e) {
      closeQuietly();
      throw new RuntimeException("Failed to initialize SharedCount", e);
    }
  }

  private String getTaskSharePath() {
    return LOCK_BASE_PATH + lockId + "/task";
  }

  private String getRecoverySharedPath() {
    return LOCK_BASE_PATH + lockId + "/recovery";
  }

  private void closeQuietly() {
    try {
      close();
    } catch (Exception e) {

View on GitHub (pinned to 86d9c8fc54)