nzw921rx opened a new issue, #12262: URL: https://github.com/apache/seatunnel/issues/12262
## Bug description SeaTunnelTransform defines open() and close() lifecycle methods, but the Flink and Spark starter adapters do not currently propagate those callbacks to the transform instance. In the Flink starter, map transforms are captured by a plain MapFunction lambda and flat-map transforms are wrapped by a plain FlatMapFunction. Flink only invokes lifecycle callbacks for RichFunction implementations, so SeaTunnelTransform.open() and SeaTunnelTransform.close() are not called. In the Spark starter, TransformMapPartitionsFunction implements FlatMapFunction and processes rows without registering task-completion cleanup. The transform is therefore not closed when a Spark task completes, fails, or is cancelled. Several transforms use lazy initialization so record processing still works without open(), but their close() methods remain unreachable on these runtimes. ## Affected code paths - Flink: seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/main/java/org/apache/seatunnel/core/starter/flink/execution/TransformExecuteProcessor.java - Spark: seatunnel-core/seatunnel-spark-starter/seatunnel-spark-starter-common/src/main/java/org/apache/seatunnel/core/starter/spark/execution/TransformExecuteProcessor.java - Lifecycle contract: seatunnel-api/src/main/java/org/apache/seatunnel/api/transform/SeaTunnelTransform.java ## Impact Resource-owning transforms can retain resources after a task or job has ended: - Python Transform can leave a Python subprocess and its stdin/stdout/stderr pump threads running in a long-lived TaskManager or Spark executor. - SQL and Calcite transforms may not close their internal SQL engines. - LLM and Embedding transforms may not close HTTP clients, connection pools, sockets, and associated threads. - Spark task retries and repeated Flink jobs can accumulate these resources and retain job classloaders. - User-defined close hooks are not executed consistently across engines. The Python Transform added by #11390 is disabled by default, but it makes the existing lifecycle gap especially visible because it owns an external operating-system process. ## Expected behavior For every runtime transform instance: 1. open() is invoked once before record processing. 2. close() is invoked exactly once when the task completes, fails, or is cancelled. 3. Cleanup attempts all owned resources even if one close operation fails. 4. Cleanup errors are logged with transform context and do not unexpectedly change an already completed job state. ## Suggested implementation ### Flink - Wrap map transforms with RichMapFunction and flat-map transforms with RichFlatMapFunction. - Forward the Flink open callback to SeaTunnelTransform.open(). - Forward the Flink close callback to SeaTunnelTransform.close(). - Keep the implementation binary-compatible with the supported Flink 1.13, 1.15, and 1.20 runtimes. ### Spark - Register one TaskContext completion listener per deserialized task function. - Initialize the transform once per task instance and close it exactly once from the completion listener. - Cover successful completion, task failure, cancellation, and retry paths. ## Test suggestions - Add starter-level lifecycle tests using a transform that records open and close counts. - Verify close is called after successful completion and after cancellation/failure on both Flink and Spark. - Add an integration assertion that a Python worker process exits after the owning task finishes. - Verify close failures do not prevent cleanup of subsequent transforms. ## Related pull request - #11390 -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
