{"record":{"id":"a52d2d243fbae693","repo":"apache/iceberg","slug":"invalid-operator-event-type-event-getclass","errorCode":null,"errorMessage":"Invalid operator event type: \" + event.getClass().getCanonicalName()","messagePattern":"Invalid operator event type: \" \\+ event\\.getClass\\(\\)\\.getCanonicalName\\(\\)","errorType":"exception","errorClass":"IllegalArgumentException","httpStatus":null,"severity":"error","filePath":"flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/DataStatisticsCoordinator.java","lineNumber":322,"sourceCode":"    }\n  }\n\n  @Override\n  public void handleEventFromOperator(int subtask, int attemptNumber, OperatorEvent event) {\n    runInCoordinatorThread(\n        () -> {\n          LOG.debug(\n              \"Handling event from subtask {} (#{}) of {}: {}\",\n              subtask,\n              attemptNumber,\n              operatorName,\n              event);\n          if (event instanceof StatisticsEvent) {\n            handleDataStatisticRequest(subtask, ((StatisticsEvent) event));\n          } else if (event instanceof RequestGlobalStatisticsEvent) {\n            handleRequestGlobalStatisticsEvent(subtask, (RequestGlobalStatisticsEvent) event);\n          } else {\n            throw new IllegalArgumentException(\n                \"Invalid operator event type: \" + event.getClass().getCanonicalName());\n          }\n        },\n        String.format(\n            Locale.ROOT,\n            \"handling operator event %s from subtask %d (#%d)\",\n            event.getClass(),\n            subtask,\n            attemptNumber));\n  }\n\n  @Override\n  public void checkpointCoordinator(long checkpointId, CompletableFuture<byte[]> resultFuture) {\n    runInCoordinatorThread(\n        () -> {\n          LOG.debug(\n              \"Snapshotting data statistics coordinator {} for checkpoint {}\",\n              operatorName,","sourceCodeStart":304,"sourceCodeEnd":340,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/shuffle/DataStatisticsCoordinator.java#L304-L340","documentation":"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.","triggerScenarios":"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.","commonSituations":"Mixed Iceberg versions in the cluster (writer operators newer than coordinator), or custom event types injected into the shuffle-statistics path.","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"],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":"// sender side: only send recognized events\nPreconditions.checkState(\n    event instanceof StatisticsEvent || event instanceof RequestGlobalStatisticsEvent,\n    \"Unsupported event for DataStatisticsCoordinator\");","typeGuard":null,"tryCatchPattern":"try {\n  coordinator.handleOperatorEvent(event);\n} catch (IllegalArgumentException e) {\n  LOG.warn(\"Dropping unknown operator event\", e);\n}","preventionTips":["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"],"tags":["flink","operator-events","shuffle-statistics"],"backgroundTag":"unexpected-response-shape","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}