{"record":{"id":"fe99a96908900c5d","repo":"apache/iceberg","slug":"invalid-operator-event-type-fe99a9","errorCode":null,"errorMessage":"Invalid operator event type: ","messagePattern":"Invalid operator event type: ","errorType":"exception","errorClass":"IllegalArgumentException","httpStatus":null,"severity":"error","filePath":"flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/TriggerManagerOperator.java","lineNumber":219,"sourceCode":"      nextEvaluationTimeState.add(nextEvaluationTime);\n    }\n\n    accumulatedChangesState.update(accumulatedChanges);\n    lastTriggerTimesState.update(lastTriggerTimes);\n    LOG.info(\n        \"Storing state: nextEvaluationTime {}, accumulatedChanges {}, lastTriggerTimes {}\",\n        nextEvaluationTime,\n        accumulatedChanges,\n        lastTriggerTimes);\n  }\n\n  @Override\n  public void handleOperatorEvent(OperatorEvent event) {\n    if (event instanceof LockReleaseEvent) {\n      LOG.info(\"Received lock released event: {}\", event);\n      handleLockRelease((LockReleaseEvent) event);\n    } else {\n      throw new IllegalArgumentException(\n          \"Invalid operator event type: \" + event.getClass().getCanonicalName());\n    }\n  }\n\n  @Override\n  public void processElement(StreamRecord<TableChange> streamRecord) throws Exception {\n    TableChange change = streamRecord.getValue();\n    accumulatedChanges.forEach(tableChange -> tableChange.merge(change));\n    if (nextEvaluationTime == null) {\n      checkAndFire(getProcessingTimeService());\n    } else {\n      LOG.info(\n          \"Trigger manager rate limiter triggered current: {}, next: {}, accumulated changes: {},{}\",\n          getProcessingTimeService().getCurrentProcessingTime(),\n          nextEvaluationTime,\n          accumulatedChanges,\n          maintenanceTaskNames);\n      rateLimiterTriggeredCounter.inc();","sourceCodeStart":201,"sourceCodeEnd":237,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/TriggerManagerOperator.java#L201-L237","documentation":"TriggerManagerOperator.handleOperatorEvent is the per-subtask OperatorEvent handler and only accepts LockReleaseEvent. Any other OperatorEvent type throws this IllegalArgumentException. The two known callers are tests (testStateRestore, testLockCheckDelay), meaning production code should only ever deliver LockReleaseEvent here.","triggerScenarios":"An OperatorEvent that is not LockReleaseEvent is delivered to the TriggerManager operator's gateway — e.g. a LockRegisterEvent sent to subtasks instead of the coordinator, or events from a different Iceberg version after restore.","commonSituations":"Custom operators/tests sending wrong event types, or cross-version savepoint restores where the event class hierarchy changed.","solutions":["Ensure LockReleaseEvent is the only event sent to the TriggerManager operator's event gateway.","Send LockRegisterEvent to TriggerManagerCoordinator (coordinator gateway), not to subtask operators.","Match Iceberg versions between checkpoint creation and restore."],"exampleFix":null,"handlingStrategy":"type-guard","validationCode":"if (!(event instanceof LockReleaseEvent)) {\n  LOG.warn(\"Ignoring unexpected event at TriggerManagerOperator: {}\", event.getClass());\n  return;\n}","typeGuard":"boolean isLockRelease(OperatorEvent e) { return e instanceof LockReleaseEvent; }","tryCatchPattern":"try {\n  operator.handleOperatorEvent(event);\n} catch (IllegalArgumentException e) {\n  LOG.error(\"Unexpected operator event type\", e);\n}","preventionTips":["Only send LockReleaseEvent to the TriggerManager operator's subtask gateway.","In tests, reuse the provided event types rather than ad-hoc OperatorEvent implementations."],"tags":["flink","operator-event","illegal-argument"],"backgroundTag":"invalid-argument-value","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-14T16:17:12.679Z"}