apache/iceberg · error · IllegalArgumentException

Invalid operator event type: " +…

Error message

Invalid operator event type: " + event.getClass().getCanonicalName()

What it means

DataStatisticsCoordinator.handleEventFromOperator processes operator events sent by subtasks and recognizes only StatisticsEvent and RequestGlobalStatisticsEvent. Any other OperatorEvent arriving at the coordinator throws an IllegalArgumentException naming the event class.

Solutions

  1. Ensure all taskmanagers/job vertices use the same Iceberg version as the coordinator
  2. Only send StatisticsEvent or RequestGlobalStatisticsEvent to this coordinator
  3. If adding new event types, extend handleEventFromOperator accordingly
Defensive patterns

Strategy: try-catch

Validate before calling

// sender side: only send recognized events
Preconditions.checkState(
    event instanceof StatisticsEvent || event instanceof RequestGlobalStatisticsEvent,
    "Unsupported event for DataStatisticsCoordinator");

Try / catch

try {
  coordinator.handleOperatorEvent(event);
} catch (IllegalArgumentException e) {
  LOG.warn("Dropping unknown operator event", e);
}

Prevention

When it happens

Trigger: A custom or unexpected OperatorEvent being sent to the coordinator, typically from modified/extended sink code or version-mismatched senders emitting new event types to an older coordinator.

Common situations: Mixed Iceberg versions in the cluster (writer operators newer than coordinator), or custom event types injected into the shuffle-statistics path.

Related errors


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

Appendix: source

Thrown at flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/DataStatisticsCoordinator.java:322

    }
  }

  @Override
  public void handleEventFromOperator(int subtask, int attemptNumber, OperatorEvent event) {
    runInCoordinatorThread(
        () -> {
          LOG.debug(
              "Handling event from subtask {} (#{}) of {}: {}",
              subtask,
              attemptNumber,
              operatorName,
              event);
          if (event instanceof StatisticsEvent) {
            handleDataStatisticRequest(subtask, ((StatisticsEvent) event));
          } else if (event instanceof RequestGlobalStatisticsEvent) {
            handleRequestGlobalStatisticsEvent(subtask, (RequestGlobalStatisticsEvent) event);
          } else {
            throw new IllegalArgumentException(
                "Invalid operator event type: " + event.getClass().getCanonicalName());
          }
        },
        String.format(
            Locale.ROOT,
            "handling operator event %s from subtask %d (#%d)",
            event.getClass(),
            subtask,
            attemptNumber));
  }

  @Override
  public void checkpointCoordinator(long checkpointId, CompletableFuture<byte[]> resultFuture) {
    runInCoordinatorThread(
        () -> {
          LOG.debug(
              "Snapshotting data statistics coordinator {} for checkpoint {}",
              operatorName,

View on GitHub (pinned to 86d9c8fc54)