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]

Reply via email to