surafel58 opened a new pull request, #11827: URL: https://github.com/apache/seatunnel/pull/11827
### Purpose of this pull request Closes #11816. `PrometheusWriter` buffered records and flushed them only on `batch_size`, on the Zeta engine timer flush (`registerFlushAction`), and in `close()`. It did not flush on checkpoint. On Spark and Flink `registerFlushAction` keeps the sink writer context's no-op default, so buffered records were held until `batch_size` was reached or the job closed. This PR overrides `prepareCommit()` to flush the buffer and return `Optional.empty()`, matching the sibling FlushSignal sinks (Doris, ClickHouse, Elasticsearch, StarRocks, MongoDB), which all flush their buffer in `prepareCommit()`. This bounds the buffered window to one checkpoint interval on every engine, without adding a connector-owned thread or any engine-level change. `flush()` throws on failure, so a failed checkpoint flush fails the checkpoint instead of silently dropping the batch. This follows up the #11778 review, where the Spark/Flink gap was flagged and documented as a limitation. The connector-level `prepareCommit` path used here is the same one the five sibling sinks already use. ### Does this PR introduce _any_ user-facing change? Yes, a behavior improvement. Before, on Spark and Flink, buffered samples were sent only when `batch_size` was reached or the writer closed, so a low-throughput streaming job could hold points in memory until it stopped. After, the buffer is also flushed on each checkpoint, so buffered samples are bounded by the checkpoint interval on all engines. The `batch_size` trigger, the Zeta timer flush, and the final flush on close are unchanged. The Prometheus sink docs and the `incompatible-changes.md` note (EN and ZH) are updated so the Spark/Flink limitation reads as checkpoint-bounded rather than no periodic flush at all. ### How was this patch tested? Added a unit test `shouldFlushOnPrepareCommitWhenEngineNeverInvokesFlushAction` in `PrometheusWriterTest`: with the engine never invoking the registered flush action (the Spark/Flink path), `prepareCommit()` flushes the buffered row and returns `Optional.empty()`. The full connector-prometheus module test suite passes locally (10 tests, 0 failures) on JDK 8. ### Check list * [x] If necessary, please update the documentation to describe the new feature. (Prometheus sink docs updated, EN and ZH) * [x] If necessary, please update `incompatible-changes.md` to describe the incompatibility caused by this PR. (updated EN and ZH to reflect checkpoint-bounded behavior) * [ ] New Jar binary package: N/A * [ ] New connector: N/A (this modifies an existing connector, no plugin-mapping / seatunnel-dist / label / plugin_config changes needed) -- 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]
