{"record":{"id":"b77f3a02762afc2f","repo":"apache/iceberg","slug":"invalid-operator-event-type-event-getclass-get-b77f3a","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/v1.20/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/v1.20/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/TriggerManagerOperator.java#L201-L237","documentation":"TriggerManagerOperator.handleOperatorEvent accepts only LockReleaseEvent (the lock-removed notification from the maintenance job). Any other OperatorEvent type is rejected with an IllegalArgumentException identifying the class. This indicates the operator received an event outside its expected protocol.","triggerScenarios":"Calling handleOperatorEvent on a TriggerManagerOperator with anything other than a LockReleaseEvent (e.g. a LockRegisterEvent or custom event). Reached in tests testStateRestore/testLockCheckDelay and via Flink's operator event pathway.","commonSituations":"Custom code or tests sending wrong event types to the trigger manager operator; wiring changes where register events are mistakenly delivered to the operator instead of the coordinator.","solutions":["Send only LockReleaseEvent to TriggerManagerOperator; LockRegisterEvent belongs on the coordinator side.","Review the event-sending code path and the OperatorEventGateway used.","If a new event type is needed, add an instanceof branch in handleOperatorEvent."],"exampleFix":"// before\noperator.handleOperatorEvent(new LockRegisterEvent(factory, id));\n// after\noperator.handleOperatorEvent(new LockReleaseEvent(taskId));","handlingStrategy":"type-guard","validationCode":"if (!(event instanceof LockReleaseEvent)) { throw new IllegalStateException(\"TriggerManagerOperator only accepts LockReleaseEvent\"); }","typeGuard":"if (event instanceof LockReleaseEvent release) { /* release handling */ }","tryCatchPattern":"try { operator.handleOperatorEvent(event); } catch (IllegalArgumentException e) { LOG.error(\"Unexpected operator event for trigger manager\", e); }","preventionTips":["Only send LockReleaseEvent to the trigger manager operator","Keep LockRegisterEvent on the coordinator side of the protocol","Extend handleOperatorEvent explicitly when adding new event types"],"tags":["flink","operator-event","iceberg-maintenance"],"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"}