SEZ9 commented on issue #12262:
URL: https://github.com/apache/seatunnel/issues/12262#issuecomment-5629337507

   Thanks for the detailed report. The analysis makes sense: Flink only invokes 
lifecycle callbacks on rich functions, so wrapping map transforms in a plain 
`MapFunction` lambda and flat-map transforms in a plain `FlatMapFunction` means 
`SeaTunnelTransform.open()`/`close()` are never called, and on Spark 
`TransformMapPartitionsFunction` has no task-completion cleanup, so `close()` 
is never reached there either.
   
   Your suggested direction looks right to me: use 
`RichMapFunction`/`RichFlatMapFunction` on Flink and forward `open`/`close`, 
and on Spark register a `TaskContext` completion listener per task function so 
`close()` runs exactly once on success, failure, cancellation, and retry. 
Keeping it binary-compatible with Flink 1.13, 1.15, and 1.20 is a good 
constraint to keep in mind. Starter-level tests with a transform that records 
open/close counts would be very helpful.
   
   Two questions before a fix PR:
   1. Have you checked whether the Zeta engine already propagates 
`open()`/`close()` for transforms, so this is Flink/Spark-only? That would help 
us document consistent lifecycle guarantees for transform authors.
   2. Since several transforms currently lazy-initialize (Python, SQL/Calcite, 
LLM, Embedding), should `open()` be made idempotent for them once it starts 
being invoked by the engines?
   
   Are you planning to open a PR for this? Happy to review.
   
   <!-- streview-comment:967 -->


-- 
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]

Reply via email to