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
- Ensure all taskmanagers/job vertices use the same Iceberg version as the coordinator
- Only send StatisticsEvent or RequestGlobalStatisticsEvent to this coordinator
- 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
- Ensure uniform Iceberg versions across all job vertices
- Do not inject custom OperatorEvents into the shuffle path
- Extend the coordinator when adding new event types
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
- Already closed files for partition:
- Altering partition keys is not supported yet.
- Altering partition keys is not supported yet.
- Altering partition keys is not supported yet.
- Altering partition keys is not supported yet.
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)