This is an automated email from the ASF dual-hosted git repository.

davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new e5f450b447 [Feature][Connector-V2][Prometheus] Migrate 
PrometheusWriter to engine-level FlushSignal (#11778)
e5f450b447 is described below

commit e5f450b44725152832404c29e1c7d551d3c26d07
Author: surafel <[email protected]>
AuthorDate: Fri Aug 14 17:54:44 2026 +0300

    [Feature][Connector-V2][Prometheus] Migrate PrometheusWriter to 
engine-level FlushSignal (#11778)
---
 docs/en/connectors/sink/Prometheus.md              |  34 +++-
 .../introduction/concepts/incompatible-changes.md  |   8 +
 docs/zh/connectors/sink/Prometheus.md              |  17 +-
 .../introduction/concepts/incompatible-changes.md  |   8 +
 .../prometheus/config/PrometheusSinkConfig.java    |   6 -
 .../prometheus/config/PrometheusSinkOptions.java   |   8 -
 .../seatunnel/prometheus/sink/PrometheusSink.java  |   2 +-
 .../prometheus/sink/PrometheusSinkFactory.java     |   1 -
 .../prometheus/sink/PrometheusWriter.java          | 172 ++++++++++-------
 .../prometheus/sink/PrometheusWriterTest.java      | 207 +++++++++++++++++++++
 .../prometheus/PrometheusTimerFlushIT.java         | 184 ++++++++++++++++++
 .../src/test/resources/prometheus_timer_flush.conf |  57 ++++++
 12 files changed, 610 insertions(+), 94 deletions(-)

diff --git a/docs/en/connectors/sink/Prometheus.md 
b/docs/en/connectors/sink/Prometheus.md
index 4bcf885b8d..f8cf6c77ed 100644
--- a/docs/en/connectors/sink/Prometheus.md
+++ b/docs/en/connectors/sink/Prometheus.md
@@ -54,7 +54,6 @@ downloaded from Maven Central.
 | retry_backoff_multiplier_ms | Int    | No       | 100     | Retry backoff 
multiplier in milliseconds. |
 | retry_backoff_max_ms        | Int    | No       | 10000   | Maximum retry 
backoff in milliseconds. |
 | batch_size                  | Int    | No       | 1024    | Maximum number 
of rows buffered before writing to Prometheus. |
-| flush_interval              | Long   | No       | 300000  | Maximum flush 
interval in milliseconds. |
 | multi_table_sink_replica    | Int    | No       | 1       | Writer replica 
count for each table in a multi-table sink job. |
 | common-options              | Config | No       | -       | Sink plugin 
common parameters. See [Sink Common 
Options](../common-options/sink-common-options.md). |
 
@@ -80,6 +79,31 @@ Supported timestamp field types:
 Replica count for multi-table sink writers. It applies to each table in a 
multi-table job. Keep the
 default value `1` unless one table needs more writer parallelism.
 
+### Timer Flush
+
+The sink can flush its buffer on a timer so that buffered samples are sent 
even when the upstream
+flow is idle and fewer than `batch_size` rows have been buffered. This timer 
is driven by the
+engine, not by the connector, and is currently supported only by **SeaTunnel 
Zeta**.
+
+Enable it by setting `sink.flush.interval` (milliseconds) in the job `env` 
block:
+
+```hocon
+env {
+  sink.flush.interval = 10000
+}
+```
+
+The engine then triggers the flush on the normal sink input-processing path, 
so there is no
+connector-owned background thread and no concurrency between the timer flush 
and the write,
+checkpoint, or close paths. A flush that fails is propagated to the engine 
instead of being silently
+dropped.
+
+> On Spark and Flink there is no periodic timer flush at all: 
`sink.flush.interval` is a Zeta engine
+> primitive, and the Spark/Flink sink writer context does not implement it. On 
those engines the
+> buffer is flushed only when it reaches `batch_size` and when the writer is 
closed. It is **not**
+> flushed on checkpoint (`PrometheusWriter` does not override 
`prepareCommit()`). For a low-throughput
+> streaming job on Spark or Flink, tune `batch_size` accordingly.
+
 ## Example
 
 ```hocon
@@ -141,15 +165,16 @@ sink {
 ## Streaming Remote Write With Batched Flush
 
 This example reads from Kafka in streaming mode and writes to a Prometheus 
remote
-write endpoint. The sink buffers up to `batch_size` rows or waits up to
-`flush_interval` milliseconds before issuing the HTTP write, which keeps 
network
-overhead low when the upstream flow is bursty.
+write endpoint. The sink buffers up to `batch_size` rows before issuing the 
HTTP
+write, and the engine-level `sink.flush.interval` (Zeta only) flushes the 
buffer
+every 10 seconds so that samples are still sent when the upstream flow is idle.
 
 ```hocon
 env {
   parallelism = 2
   job.mode = "STREAMING"
   checkpoint.interval = 30000
+  sink.flush.interval = 10000
 }
 
 source {
@@ -176,7 +201,6 @@ sink {
     key_value = "c_double"
     key_timestamp = "c_timestamp"
     batch_size = 2048
-    flush_interval = 10000
     retry = 5
     retry_backoff_multiplier_ms = 200
     retry_backoff_max_ms = 10000
diff --git a/docs/en/introduction/concepts/incompatible-changes.md 
b/docs/en/introduction/concepts/incompatible-changes.md
index 7ad6a2af1b..8f15b63ce1 100644
--- a/docs/en/introduction/concepts/incompatible-changes.md
+++ b/docs/en/introduction/concepts/incompatible-changes.md
@@ -120,6 +120,14 @@ You need to check this document before you upgrade to 
related version.
   - **Impact**: Existing jobs that read POI-engine Excel files larger than 50 
MB - which previously succeeded at the cost of heavy memory pressure - will now 
fail fast with a `FileConnectorException` instead of potentially OOMing the 
worker.
   - **Migration Guide**: For POI jobs that must read large Excel files and 
have sufficient worker memory, raise the limit with `poi_excel_max_file_size = 
<bytes>`. Otherwise switch to `excel_engine = EasyExcel`, which streams rows 
lazily and is not subject to the limit.
 
+- **Breaking Change: Prometheus Sink `flush_interval` option removed**
+  - **Affected component**: `seatunnel-connectors-v2/connector-prometheus`
+  - **Description**: The Prometheus Sink no longer starts its own background 
flush thread. The connector-level `flush_interval` option has been removed. 
Timer-based flushing is now driven by the engine through `sink.flush.interval` 
in the job `env` block, which is **supported only by the Zeta engine**.
+  - **Impact**:
+    - **Spark and Flink lose periodic timer-based flushing.** The removed 
`flush_interval` scheduler was a plain connector-owned thread that ran on all 
engines. Its replacement, `sink.flush.interval`, is a Zeta engine primitive; 
the Spark and Flink sink writer contexts do not implement it, so there is no 
periodic flush on those engines. On Spark and Flink the buffer is now flushed 
only when it reaches `batch_size` and when the writer is closed (it is not 
flushed on checkpoint). A low-thr [...]
+    - A leftover `flush_interval` key in the `Prometheus` sink block is 
rejected only when the config is validated with `--check` / `--dry-run=static` 
/ `--dry-run=connect` (which run `validateUnknownKeys`). A directly submitted 
job silently ignores the stray key; the connector logs a warning once per sink 
writer at startup instead (so a job with parallelism N, multiple tables, or 
replicas logs it multiple times).
+  - **Migration Guide**: Remove `flush_interval` from the `Prometheus` sink 
block. To keep timer-based flushing on Zeta, set `sink.flush.interval` 
(milliseconds) in the job `env` block. On Spark and Flink, rely on 
`batch_size`. The `batch_size` trigger and the final flush on writer close are 
unchanged on all engines.
+
 ### Transform Changes
 
 - **[BREAKING]** SQL Transform `PARSEDATETIME`, `TO_DATE`, and `IS_DATE` 
functions now only accept whitelisted datetime format patterns. Custom format 
patterns that were previously accepted will now fail at runtime. The supported 
patterns are:
diff --git a/docs/zh/connectors/sink/Prometheus.md 
b/docs/zh/connectors/sink/Prometheus.md
index 7925219f7b..558f0ca0db 100644
--- a/docs/zh/connectors/sink/Prometheus.md
+++ b/docs/zh/connectors/sink/Prometheus.md
@@ -51,7 +51,6 @@ Prometheus 数据接收器把上游数据写入 Prometheus remote write API。
 | retry_backoff_multiplier_ms | Int    | 否       | 100    | 重试退避时间倍数,单位毫秒。 |
 | retry_backoff_max_ms        | Int    | 否       | 10000  | 最大重试退避时间,单位毫秒。 |
 | batch_size                  | Int    | 否       | 1024   | 写入 Prometheus 
前最多缓存的行数。 |
-| flush_interval              | Long   | 否       | 300000 | 最大刷新间隔,单位毫秒。 |
 | multi_table_sink_replica    | Int    | 否       | 1      | 
多表写入时,每张表使用的写入器副本数。 |
 | common-options              | Config | 否       | -      | 
接收器插件通用参数,详情请参考[接收器通用选项](../common-options/sink-common-options.md)。 |
 
@@ -75,6 +74,22 @@ Sink 会自动补充 remote write 需要的请求头:`Content-type`、`Content
 
 多表写入时,每张表使用的 Sink Writer 副本数。默认值为 `1`;只有当单张表需要更高写入并行度时才建议调大。
 
+### 定时刷新
+
+即使上游数据空闲、缓存的行数还没达到 
`batch_size`,接收器也可以按定时器刷新缓存,把已缓存的采样点发送出去。该定时器由引擎驱动,而不是由连接器驱动,**目前仅 SeaTunnel 
Zeta 支持**。
+
+在作业的 `env` 中设置 `sink.flush.interval`(单位毫秒)即可启用:
+
+```hocon
+env {
+  sink.flush.interval = 10000
+}
+```
+
+引擎会在正常的 Sink 
数据处理线程上触发刷新,因此不需要连接器自己维护后台线程,也不会和写入、检查点、关闭等流程产生并发。刷新失败会被抛给引擎,而不会被静默丢弃。
+
+> 在 Spark 和 Flink 上完全没有周期性定时刷新:`sink.flush.interval` 是 Zeta 引擎的能力,Spark/Flink 
的 Sink 写入器上下文并未实现它。在这两个引擎上,缓存只会在达到 `batch_size` 
以及写入器关闭时被刷新,**不会**在检查点时刷新(`PrometheusWriter` 未重写 
`prepareCommit()`)。对于低吞吐的流式作业,请相应调整 `batch_size`。
+
 ## 示例
 
 ```hocon
diff --git a/docs/zh/introduction/concepts/incompatible-changes.md 
b/docs/zh/introduction/concepts/incompatible-changes.md
index a44dc4e8e0..a005482e49 100644
--- a/docs/zh/introduction/concepts/incompatible-changes.md
+++ b/docs/zh/introduction/concepts/incompatible-changes.md
@@ -115,6 +115,14 @@
   - **影响**:此前以 POI 引擎读取大于 50 MB Excel 文件的任务(虽然成功但伴随严重内存压力)现在会以 
`FileConnectorException` 快速失败,而不再可能导致 worker OOM。
   - **迁移指南**:对于必须读取大 Excel 文件且 worker 内存充足的 POI 任务,可通过 
`poi_excel_max_file_size = <字节数>` 调高限制;否则切换为 `excel_engine = 
EasyExcel`,该引擎惰性流式读取行,不受此限制约束。
 
+- **破坏性变更:移除 Prometheus Sink 的 `flush_interval` 选项**
+  - **受影响组件**:`seatunnel-connectors-v2/connector-prometheus`
+  - **变更说明**:Prometheus Sink 不再启动自己的后台刷新线程,连接器级的 `flush_interval` 
选项已被移除。定时刷新改为由引擎通过作业 `env` 中的 `sink.flush.interval` 驱动,**仅 Zeta 引擎支持**。
+  - **影响**:
+    - **Spark 和 Flink 会失去周期性定时刷新。** 被移除的 `flush_interval` 
调度器是连接器自己的线程,在所有引擎上都能工作;其替代者 `sink.flush.interval` 是 Zeta 引擎的能力,Spark 和 Flink 的 
Sink 写入器上下文并未实现它,因此这两个引擎上没有周期性刷新。在 Spark 和 Flink 上,缓存现在只会在达到 `batch_size` 
以及写入器关闭时被刷新(不会在检查点时刷新)。因此低吞吐的流式作业可能会把缓存的采样点一直保存在内存中直到作业停止;请相应调整 `batch_size`。
+    - 只有在使用 `--check` / `--dry-run=static` / `--dry-run=connect` 校验配置时(会执行 
`validateUnknownKeys`),`Prometheus` sink 中残留的 `flush_interval` 
键才会被拒绝。直接提交的作业会静默忽略该残留键;连接器会在每个 Sink 写入器启动时各打印一次告警作为替代提示(因此并行度为 
N、多表或多副本的作业会多次打印)。
+  - **迁移指南**:从 `Prometheus` sink 中移除 `flush_interval`。如需在 Zeta 上继续使用定时刷新,请在作业 
`env` 中设置 `sink.flush.interval`(毫秒)。在 Spark 和 Flink 上请依赖 
`batch_size`。`batch_size` 触发和写入器关闭时的最后一次刷新在所有引擎上保持不变。
+
 ### 转换变更
 
 - **[BREAKING]** SQL Transform 的 `PARSEDATETIME`、`TO_DATE` 和 `IS_DATE` 
函数现在只接受白名单中的日期时间格式模式。以前接受的自定义格式模式现在将在运行时失败。支持的模式有:
diff --git 
a/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/config/PrometheusSinkConfig.java
 
b/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/config/PrometheusSinkConfig.java
index 1505e37647..f944ffc9fb 100644
--- 
a/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/config/PrometheusSinkConfig.java
+++ 
b/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/config/PrometheusSinkConfig.java
@@ -38,8 +38,6 @@ public class PrometheusSinkConfig extends HttpConfig {
 
     private int batchSize;
 
-    private long flushInterval;
-
     public static PrometheusSinkConfig loadConfig(ReadonlyConfig pluginConfig) 
{
         PrometheusSinkConfig sinkConfig = new PrometheusSinkConfig();
         if 
(pluginConfig.getOptional(PrometheusSinkOptions.KEY_VALUE).isPresent()) {
@@ -56,10 +54,6 @@ public class PrometheusSinkConfig extends HttpConfig {
         // would leave batchSize at 0 and disable the size-based flush trigger.
         int batchSize = 
checkIntArgument(pluginConfig.get(PrometheusSinkOptions.BATCH_SIZE));
         sinkConfig.setBatchSize(batchSize);
-        if 
(pluginConfig.getOptional(PrometheusSinkOptions.FLUSH_INTERVAL).isPresent()) {
-            long flushInterval = 
pluginConfig.get(PrometheusSinkOptions.FLUSH_INTERVAL);
-            sinkConfig.setFlushInterval(flushInterval);
-        }
         return sinkConfig;
     }
 
diff --git 
a/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/config/PrometheusSinkOptions.java
 
b/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/config/PrometheusSinkOptions.java
index 0ebbdffc3b..427994f765 100644
--- 
a/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/config/PrometheusSinkOptions.java
+++ 
b/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/config/PrometheusSinkOptions.java
@@ -25,8 +25,6 @@ public class PrometheusSinkOptions extends HttpCommonOptions {
 
     private static final int DEFAULT_BATCH_SIZE = 1024;
 
-    private static final Long DEFAULT_FLUSH_INTERVAL = 300000L;
-
     public static final Option<String> KEY_TIMESTAMP =
             Options.key("key_timestamp")
                     .stringType()
@@ -44,10 +42,4 @@ public class PrometheusSinkOptions extends HttpCommonOptions 
{
                     .intType()
                     .defaultValue(DEFAULT_BATCH_SIZE)
                     .withDescription("the batch size writer to prometheus");
-
-    public static final Option<Long> FLUSH_INTERVAL =
-            Options.key("flush_interval")
-                    .longType()
-                    .defaultValue(DEFAULT_FLUSH_INTERVAL)
-                    .withDescription("the flush interval writer to 
prometheus");
 }
diff --git 
a/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusSink.java
 
b/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusSink.java
index 15681f429c..533364f057 100644
--- 
a/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusSink.java
+++ 
b/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusSink.java
@@ -69,7 +69,7 @@ public class PrometheusSink extends 
AbstractSimpleSink<SeaTunnelRow, Void>
     @Override
     public PrometheusWriter createWriter(SinkWriter.Context context) {
         return new PrometheusWriter(
-                catalogTable.getSeaTunnelRowType(), httpParameter, 
pluginConfig);
+                catalogTable.getSeaTunnelRowType(), httpParameter, 
pluginConfig, context);
     }
 
     @Override
diff --git 
a/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusSinkFactory.java
 
b/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusSinkFactory.java
index 35c2a9ad36..879c6a2862 100644
--- 
a/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusSinkFactory.java
+++ 
b/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusSinkFactory.java
@@ -53,7 +53,6 @@ public class PrometheusSinkFactory extends HttpSinkFactory {
                 .optional(PrometheusSinkOptions.RETRY_BACKOFF_MULTIPLIER_MS)
                 .optional(PrometheusSinkOptions.RETRY_BACKOFF_MAX_MS)
                 .optional(PrometheusSinkOptions.BATCH_SIZE)
-                .optional(PrometheusSinkOptions.FLUSH_INTERVAL)
                 .optional(SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA)
                 .build();
     }
diff --git 
a/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusWriter.java
 
b/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusWriter.java
index 307abb8eed..f1de9b2c91 100644
--- 
a/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusWriter.java
+++ 
b/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusWriter.java
@@ -17,6 +17,7 @@
 package org.apache.seatunnel.connectors.seatunnel.prometheus.sink;
 
 import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.sink.SinkWriter;
 import org.apache.seatunnel.api.table.type.SeaTunnelRow;
 import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
 import org.apache.seatunnel.common.exception.CommonErrorCodeDeprecated;
@@ -42,33 +43,30 @@ import java.io.IOException;
 import java.util.ArrayList;
 import java.util.List;
 import java.util.Map;
-import java.util.concurrent.Executors;
-import java.util.concurrent.ScheduledExecutorService;
-import java.util.concurrent.ScheduledFuture;
-import java.util.concurrent.TimeUnit;
 
 @Slf4j
 public class PrometheusWriter extends HttpSinkWriter {
+
+    // The removed connector-level option key, kept only to detect and warn 
about a leftover key in
+    // an upgraded job config.
+    private static final String REMOVED_FLUSH_INTERVAL_KEY = "flush_interval";
+
     private final List<Point> batchList;
-    private volatile Exception flushException;
     private final Integer batchSize;
-    private final long flushInterval;
-    private PrometheusSinkConfig sinkConfig;
+    private final PrometheusSinkConfig sinkConfig;
     private final Serializer serializer;
     protected final HttpClientProvider httpClient;
-    private ScheduledExecutorService executor;
-    private ScheduledFuture scheduledFuture;
 
     public PrometheusWriter(
             SeaTunnelRowType seaTunnelRowType,
             HttpParameter httpParameter,
-            ReadonlyConfig pluginConfig) {
+            ReadonlyConfig pluginConfig,
+            SinkWriter.Context context) {
 
         super(seaTunnelRowType, httpParameter);
         this.batchList = new ArrayList<>();
         this.sinkConfig = PrometheusSinkConfig.loadConfig(pluginConfig);
         this.batchSize = sinkConfig.getBatchSize();
-        this.flushInterval = sinkConfig.getFlushInterval();
         this.serializer =
                 new PrometheusSerializer(
                         seaTunnelRowType,
@@ -76,23 +74,27 @@ public class PrometheusWriter extends HttpSinkWriter {
                         sinkConfig.getKeyLabel(),
                         sinkConfig.getKeyValue());
         this.httpClient = new HttpClientProvider(httpParameter);
-        if (flushInterval > 0) {
-            log.info("start schedule submit message,interval:{}", 
flushInterval);
-            this.executor =
-                    Executors.newScheduledThreadPool(
-                            1,
-                            runnable -> {
-                                Thread thread = new Thread(runnable);
-                                thread.setDaemon(true);
-                                thread.setName("Prometheus-Metric-Sender");
-                                return thread;
-                            });
-            this.scheduledFuture =
-                    executor.scheduleAtFixedRate(
-                            this::flushSchedule,
-                            flushInterval,
-                            flushInterval,
-                            TimeUnit.MILLISECONDS);
+        // The connector-level `flush_interval` option was removed in favor of 
the engine-level
+        // `sink.flush.interval`. A leftover key in an upgraded job config is 
silently ignored on a
+        // direct job run (only `--check`/`--dry-run` reject unknown keys), so 
warn here (once per
+        // writer instance) to give operators a signal instead of silently 
dropping periodic
+        // flushing.
+        if 
(pluginConfig.getSourceMap().containsKey(REMOVED_FLUSH_INTERVAL_KEY)) {
+            log.warn(
+                    "The connector option 'flush_interval' has been removed 
and is ignored. Use the "
+                            + "engine-level 'sink.flush.interval' in the job 
'env' block instead. "
+                            + "Engine-level timer flush is supported only by 
Zeta; on Spark and "
+                            + "Flink there is no periodic flush, so tune 
'batch_size' instead.");
+        }
+        // Opt in to engine-level timer flush. On Zeta the engine invokes this 
action on the normal
+        // Sink input-processing path when a FlushSignal arrives, so there is 
no connector-owned
+        // scheduler thread and no concurrency with write/checkpoint/close. On 
Spark and Flink the
+        // Context does not implement registerFlushAction (it keeps the 
interface's no-op default),
+        // so there is no periodic timer flush there; the buffer is flushed on 
batch_size and on
+        // close(). The null-check is defensive for non-standard/test call 
sites that may not supply
+        // a context.
+        if (context != null) {
+            context.registerFlushAction(this::flush);
         }
     }
 
@@ -103,8 +105,6 @@ public class PrometheusWriter extends HttpSinkWriter {
     }
 
     public void write(Point record) {
-        checkFlushException();
-
         synchronized (batchList) {
             batchList.add(record);
             if (batchSize > 0 && batchList.size() >= batchSize) {
@@ -113,45 +113,38 @@ public class PrometheusWriter extends HttpSinkWriter {
         }
     }
 
-    private void flushSchedule() {
-        synchronized (batchList) {
-            if (!batchList.isEmpty()) {
-                flush();
-            }
-        }
-    }
-
-    private void checkFlushException() {
-        if (flushException != null) {
-            throw new PrometheusConnectorException(
-                    CommonErrorCodeDeprecated.FLUSH_DATA_FAILED,
-                    "Writing records to prometheus failed.",
-                    flushException);
-        }
-    }
-
     private void flush() {
-        checkFlushException();
-        if (batchList.isEmpty()) {
-            return;
-        }
-        try {
-            byte[] body = snappy(batchList);
-            ByteArrayEntity byteArrayEntity = new ByteArrayEntity(body);
-            HttpResponse response =
-                    httpClient.doPost(
-                            httpParameter.getUrl(), 
httpParameter.getHeaders(), byteArrayEntity);
-            if (HttpStatus.SC_NO_CONTENT == response.getCode()) {
+        synchronized (batchList) {
+            if (batchList.isEmpty()) {
                 return;
             }
-            log.error(
-                    "http client execute exception, http response status 
code:[{}], content:[{}]",
-                    response.getCode(),
-                    response.getContent());
-        } catch (Exception e) {
-            log.error(e.getMessage(), e);
-        } finally {
-            batchList.clear();
+            try {
+                byte[] body = snappy(batchList);
+                ByteArrayEntity byteArrayEntity = new ByteArrayEntity(body);
+                HttpResponse response =
+                        httpClient.doPost(
+                                httpParameter.getUrl(),
+                                httpParameter.getHeaders(),
+                                byteArrayEntity);
+                if (HttpStatus.SC_NO_CONTENT == response.getCode()) {
+                    batchList.clear();
+                    return;
+                }
+                // Propagate the failure to the engine instead of silently 
dropping the batch, so a
+                // flush that did not succeed is not treated as a successful 
flush.
+                throw new PrometheusConnectorException(
+                        CommonErrorCodeDeprecated.FLUSH_DATA_FAILED,
+                        String.format(
+                                "Writing records to prometheus failed, http 
response status code:[%d], content:[%s]",
+                                response.getCode(), response.getContent()));
+            } catch (PrometheusConnectorException e) {
+                throw e;
+            } catch (Exception e) {
+                throw new PrometheusConnectorException(
+                        CommonErrorCodeDeprecated.FLUSH_DATA_FAILED,
+                        "Writing records to prometheus failed.",
+                        e);
+            }
         }
     }
 
@@ -202,13 +195,48 @@ public class PrometheusWriter extends HttpSinkWriter {
 
     @Override
     public void close() throws IOException {
-        super.close();
-        if (scheduledFuture != null) {
-            scheduledFuture.cancel(false);
-            if (executor != null) {
-                executor.shutdownNow();
+        // Run the final flush and both cleanup steps unconditionally, but 
keep the first failure as
+        // the primary exception and attach later ones with addSuppressed. 
Otherwise an IOException
+        // from closing an HTTP client (thrown from a finally block) would 
replace the meaningful
+        // "Writing records to prometheus failed" exception from the final 
flush.
+        Throwable primary = null;
+        try {
+            // Send any records still buffered before the writer is closed.
+            flush();
+        } catch (Throwable t) {
+            primary = t;
+        }
+        try {
+            // Close the HttpClientProvider actually used for remote-write 
(this field shadows the
+            // parent's), otherwise it would leak when the writer is closed.
+            httpClient.close();
+        } catch (Throwable t) {
+            primary = addAsPrimaryOrSuppressed(primary, t);
+        }
+        try {
+            super.close();
+        } catch (Throwable t) {
+            primary = addAsPrimaryOrSuppressed(primary, t);
+        }
+        if (primary != null) {
+            if (primary instanceof IOException) {
+                throw (IOException) primary;
+            }
+            if (primary instanceof RuntimeException) {
+                throw (RuntimeException) primary;
+            }
+            if (primary instanceof Error) {
+                throw (Error) primary;
             }
+            throw new IOException(primary);
+        }
+    }
+
+    private static Throwable addAsPrimaryOrSuppressed(Throwable primary, 
Throwable next) {
+        if (primary == null) {
+            return next;
         }
-        this.flush();
+        primary.addSuppressed(next);
+        return primary;
     }
 }
diff --git 
a/seatunnel-connectors-v2/connector-prometheus/src/test/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusWriterTest.java
 
b/seatunnel-connectors-v2/connector-prometheus/src/test/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusWriterTest.java
new file mode 100644
index 0000000000..ac769bf96f
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-prometheus/src/test/java/org/apache/seatunnel/connectors/seatunnel/prometheus/sink/PrometheusWriterTest.java
@@ -0,0 +1,207 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.prometheus.sink;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.sink.SinkWriter;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.common.utils.function.RunnableWithException;
+import 
org.apache.seatunnel.connectors.seatunnel.http.client.HttpClientProvider;
+import org.apache.seatunnel.connectors.seatunnel.http.client.HttpResponse;
+import org.apache.seatunnel.connectors.seatunnel.http.config.HttpParameter;
+import 
org.apache.seatunnel.connectors.seatunnel.prometheus.Exception.PrometheusConnectorException;
+
+import org.apache.http.HttpStatus;
+import org.apache.http.entity.ByteArrayEntity;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.MockedConstruction;
+
+import java.io.IOException;
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+class PrometheusWriterTest {
+
+    /**
+     * The writer should opt in to engine-level timer flush by registering a 
flush action, and that
+     * action should send the buffered records only when the engine invokes it 
(the flush signal),
+     * not on every write.
+     */
+    @Test
+    void shouldRegisterFlushActionAndFlushBufferedRecordsOnSignal() throws 
Exception {
+        SinkWriter.Context context = mock(SinkWriter.Context.class);
+        HttpResponse ok = new HttpResponse(HttpStatus.SC_NO_CONTENT);
+        ArgumentCaptor<RunnableWithException> actionCaptor =
+                ArgumentCaptor.forClass(RunnableWithException.class);
+
+        try (MockedConstruction<HttpClientProvider> ignored =
+                mockConstruction(
+                        HttpClientProvider.class,
+                        (mockClient, ctx) ->
+                                when(mockClient.doPost(
+                                                anyString(), any(), 
any(ByteArrayEntity.class)))
+                                        .thenReturn(ok))) {
+
+            PrometheusWriter writer = createWriter(context);
+            writer.write(newPoint());
+
+            verify(context, 
times(1)).registerFlushAction(actionCaptor.capture());
+            // Buffered only: nothing is sent until the engine delivers a 
flush signal.
+            verify(writer.httpClient, never())
+                    .doPost(anyString(), any(), any(ByteArrayEntity.class));
+
+            actionCaptor.getValue().run();
+
+            verify(writer.httpClient, times(1))
+                    .doPost(anyString(), any(), any(ByteArrayEntity.class));
+        }
+    }
+
+    /**
+     * A flush that does not succeed must be propagated to the engine instead 
of being silently
+     * treated as a successful flush.
+     */
+    @Test
+    void shouldPropagateFlushFailure() throws Exception {
+        SinkWriter.Context context = mock(SinkWriter.Context.class);
+        HttpResponse failed = new HttpResponse(HttpStatus.SC_BAD_REQUEST, 
"boom");
+        ArgumentCaptor<RunnableWithException> actionCaptor =
+                ArgumentCaptor.forClass(RunnableWithException.class);
+
+        try (MockedConstruction<HttpClientProvider> ignored =
+                mockConstruction(
+                        HttpClientProvider.class,
+                        (mockClient, ctx) ->
+                                when(mockClient.doPost(
+                                                anyString(), any(), 
any(ByteArrayEntity.class)))
+                                        .thenReturn(failed))) {
+
+            PrometheusWriter writer = createWriter(context);
+            writer.write(newPoint());
+
+            verify(context, 
times(1)).registerFlushAction(actionCaptor.capture());
+            Assertions.assertThrows(
+                    PrometheusConnectorException.class, () -> 
actionCaptor.getValue().run());
+        }
+    }
+
+    /**
+     * On Spark and Flink the sink writer context does not implement 
registerFlushAction (it keeps
+     * the interface's no-op default), so the engine never invokes the flush 
action. The buffered
+     * records must still be delivered when the writer is closed, not lost.
+     */
+    @Test
+    void shouldFlushOnCloseWhenEngineNeverInvokesFlushAction() throws 
Exception {
+        SinkWriter.Context context = mock(SinkWriter.Context.class);
+        HttpResponse ok = new HttpResponse(HttpStatus.SC_NO_CONTENT);
+
+        try (MockedConstruction<HttpClientProvider> ignored =
+                mockConstruction(
+                        HttpClientProvider.class,
+                        (mockClient, ctx) ->
+                                when(mockClient.doPost(
+                                                anyString(), any(), 
any(ByteArrayEntity.class)))
+                                        .thenReturn(ok))) {
+
+            PrometheusWriter writer = createWriter(context);
+            writer.write(newPoint());
+
+            // Simulate Spark/Flink: the registered flush action is never 
invoked by the engine.
+            verify(writer.httpClient, never())
+                    .doPost(anyString(), any(), any(ByteArrayEntity.class));
+
+            writer.close();
+
+            // close() must flush the buffered row so it is not lost.
+            verify(writer.httpClient, times(1))
+                    .doPost(anyString(), any(), any(ByteArrayEntity.class));
+        }
+    }
+
+    /**
+     * When the final flush in close() fails and closing the HTTP client also 
throws, the meaningful
+     * flush failure must be the exception that surfaces, with the 
client-close error suppressed,
+     * not the other way around.
+     */
+    @Test
+    void closeShouldKeepFlushExceptionWhenHttpClientCloseAlsoThrows() throws 
Exception {
+        SinkWriter.Context context = mock(SinkWriter.Context.class);
+        HttpResponse failed = new HttpResponse(HttpStatus.SC_BAD_REQUEST, 
"boom");
+
+        try (MockedConstruction<HttpClientProvider> ignored =
+                mockConstruction(
+                        HttpClientProvider.class,
+                        (mockClient, ctx) ->
+                                when(mockClient.doPost(
+                                                anyString(), any(), 
any(ByteArrayEntity.class)))
+                                        .thenReturn(failed))) {
+
+            PrometheusWriter writer = createWriter(context);
+            writer.write(newPoint());
+            // The HTTP client teardown also fails during close().
+            doThrow(new IOException("client close 
failed")).when(writer.httpClient).close();
+
+            PrometheusConnectorException thrown =
+                    
Assertions.assertThrows(PrometheusConnectorException.class, writer::close);
+
+            boolean clientCloseSuppressed = false;
+            for (Throwable suppressed : thrown.getSuppressed()) {
+                if (suppressed instanceof IOException) {
+                    clientCloseSuppressed = true;
+                }
+            }
+            Assertions.assertTrue(
+                    clientCloseSuppressed,
+                    "The client-close IOException should be suppressed on the 
flush failure");
+        }
+    }
+
+    private PrometheusWriter createWriter(SinkWriter.Context context) {
+        HttpParameter httpParameter = new HttpParameter();
+        httpParameter.setUrl("http://localhost:9090/api/v1/write";);
+        httpParameter.setHeaders(new HashMap<>());
+        SeaTunnelRowType rowType =
+                new SeaTunnelRowType(
+                        new String[] {"value"}, new SeaTunnelDataType[] 
{BasicType.DOUBLE_TYPE});
+        // batch_size is not set, so it resolves to its declared default of 
1024; each test writes a
+        // single row, which stays well below that, so no size-triggered flush 
occurs.
+        ReadonlyConfig pluginConfig = ReadonlyConfig.fromMap(new HashMap<>());
+        return new PrometheusWriter(rowType, httpParameter, pluginConfig, 
context);
+    }
+
+    private Point newPoint() {
+        Map<String, String> metric = new HashMap<>();
+        metric.put("__name__", "test_metric");
+        return 
Point.builder().metric(metric).value(1.0).timestamp(1_600_000_000_000L).build();
+    }
+}
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-prometheus-e2e/src/test/java/org/apache/seatunnel/e2e/connector/prometheus/PrometheusTimerFlushIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-prometheus-e2e/src/test/java/org/apache/seatunnel/e2e/connector/prometheus/PrometheusTimerFlushIT.java
new file mode 100644
index 0000000000..30b1a866dc
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-prometheus-e2e/src/test/java/org/apache/seatunnel/e2e/connector/prometheus/PrometheusTimerFlushIT.java
@@ -0,0 +1,184 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.seatunnel.e2e.connector.prometheus;
+
+import org.apache.seatunnel.common.utils.JsonUtils;
+import org.apache.seatunnel.e2e.common.TestResource;
+import org.apache.seatunnel.e2e.common.TestSuiteBase;
+import org.apache.seatunnel.e2e.common.container.EngineType;
+import org.apache.seatunnel.e2e.common.container.TestContainer;
+import org.apache.seatunnel.e2e.common.junit.DisabledOnContainer;
+import org.apache.seatunnel.e2e.common.util.JobIdGenerator;
+
+import org.apache.http.client.methods.CloseableHttpResponse;
+import org.apache.http.client.methods.HttpGet;
+import org.apache.http.impl.client.CloseableHttpClient;
+import org.apache.http.impl.client.HttpClients;
+import org.apache.http.util.EntityUtils;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.TestTemplate;
+import org.testcontainers.containers.Container;
+import org.testcontainers.containers.GenericContainer;
+import org.testcontainers.containers.output.Slf4jLogConsumer;
+import org.testcontainers.containers.wait.strategy.HostPortWaitStrategy;
+import org.testcontainers.lifecycle.Startables;
+import org.testcontainers.utility.DockerImageName;
+import org.testcontainers.utility.DockerLoggerFactory;
+
+import com.jayway.jsonpath.JsonPath;
+import lombok.Data;
+import lombok.extern.slf4j.Slf4j;
+
+import java.time.Duration;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Stream;
+
+import static org.awaitility.Awaitility.await;
+
+@Slf4j
+@DisabledOnContainer(
+        value = {},
+        type = {EngineType.SPARK, EngineType.FLINK},
+        disabledReason =
+                "engine-level timer flush (sink.flush.interval) is only 
supported on Zeta engine")
+public class PrometheusTimerFlushIT extends TestSuiteBase implements 
TestResource {
+
+    private static final String IMAGE = "bitnamilegacy/prometheus:2.53.0";
+
+    private static final String HOST = "prometheus-host";
+
+    private static final String METRIC_NAME = "timer_flush_metric";
+
+    private GenericContainer<?> prometheusContainer;
+
+    @BeforeAll
+    @Override
+    public void startUp() {
+        this.prometheusContainer =
+                new GenericContainer<>(DockerImageName.parse(IMAGE))
+                        .withNetwork(NETWORK)
+                        .withNetworkAliases(HOST)
+                        .withEnv("TZ", "Asia/Shanghai")
+                        .withExposedPorts(9090)
+                        .withCommand(
+                                
"--config.file=/opt/bitnami/prometheus/conf/prometheus.yml",
+                                "--web.enable-remote-write-receiver")
+                        .withLogConsumer(new 
Slf4jLogConsumer(DockerLoggerFactory.getLogger(IMAGE)))
+                        .waitingFor(
+                                new HostPortWaitStrategy()
+                                        
.withStartupTimeout(Duration.ofMinutes(2)));
+        Startables.deepStart(Stream.of(prometheusContainer)).join();
+        log.info("Prometheus container started");
+    }
+
+    @AfterAll
+    @Override
+    public void tearDown() {
+        if (prometheusContainer != null) {
+            prometheusContainer.stop();
+        }
+    }
+
+    @TestTemplate
+    public void testPrometheusTimerFlush(TestContainer container) throws 
Exception {
+        String jobId = String.valueOf(JobIdGenerator.newJobId());
+        CompletableFuture<Container.ExecResult> jobFuture =
+                CompletableFuture.supplyAsync(
+                        () -> {
+                            try {
+                                return 
container.executeJob("/prometheus_timer_flush.conf", jobId);
+                            } catch (Exception e) {
+                                throw new RuntimeException(e);
+                            }
+                        });
+
+        try {
+            // Wait until the streaming job is actually running.
+            await().atMost(2, TimeUnit.MINUTES)
+                    .pollInterval(2, TimeUnit.SECONDS)
+                    .untilAsserted(
+                            () -> {
+                                if (jobFuture.isDone()) {
+                                    Container.ExecResult jobResult = 
jobFuture.get();
+                                    Assertions.fail(
+                                            "The streaming job terminated 
before reaching RUNNING: "
+                                                    + jobResult.getStderr());
+                                }
+                                Assertions.assertEquals("RUNNING", 
container.getJobStatus(jobId));
+                            });
+
+            // batch_size (100) is larger than the single buffered row, and 
the checkpoint
+            // interval is long, so the row can only reach Prometheus through 
the engine timer
+            // flush. Assert it arrives while the job is still running (before 
the writer closes).
+            await().atMost(120, TimeUnit.SECONDS)
+                    .ignoreExceptions()
+                    .pollInterval(2, TimeUnit.SECONDS)
+                    .untilAsserted(
+                            () -> {
+                                Assertions.assertFalse(
+                                        jobFuture.isDone(),
+                                        "The streaming job must still be 
running when timer flush publishes the buffered row");
+                                Metric metric = queryMetric(METRIC_NAME);
+                                Assertions.assertNotNull(
+                                        metric, "Prometheus has not received 
the buffered row yet");
+                                Assertions.assertEquals(
+                                        METRIC_NAME, 
metric.getMetric().get("__name__"));
+                                Assertions.assertEquals("2.34", 
metric.getValue().get(1));
+                            });
+        } finally {
+            if (!jobFuture.isDone()) {
+                Container.ExecResult cancelResult = container.cancelJob(jobId);
+                Assertions.assertEquals(0, cancelResult.getExitCode(), 
cancelResult.getStderr());
+            }
+        }
+    }
+
+    private Metric queryMetric(String metricName) throws Exception {
+        try (CloseableHttpClient httpClient = HttpClients.createDefault()) {
+            HttpGet httpGet =
+                    new HttpGet(
+                            "http://";
+                                    + prometheusContainer.getHost()
+                                    + ":"
+                                    + prometheusContainer.getMappedPort(9090)
+                                    + "/api/v1/query?query="
+                                    + metricName);
+            try (CloseableHttpResponse response = httpClient.execute(httpGet)) 
{
+                String responseContent = 
EntityUtils.toString(response.getEntity());
+                List<Metric> metrics =
+                        JsonUtils.toList(
+                                JsonPath.read(responseContent, 
"$.data.result.*").toString(),
+                                Metric.class);
+                return metrics.isEmpty() ? null : metrics.get(0);
+            }
+        }
+    }
+
+    @Data
+    public static class Metric {
+
+        private Map<String, String> metric;
+
+        private List<String> value;
+    }
+}
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-prometheus-e2e/src/test/resources/prometheus_timer_flush.conf
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-prometheus-e2e/src/test/resources/prometheus_timer_flush.conf
new file mode 100644
index 0000000000..baef0c9031
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-prometheus-e2e/src/test/resources/prometheus_timer_flush.conf
@@ -0,0 +1,57 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements.  See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License.  You may obtain a copy of the License at
+#
+#    http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+# Streaming job that verifies engine-level timer flush (sink.flush.interval) 
for the Prometheus sink.
+# The FakeSource emits one metric row and then stays open (streaming source 
never signals end), so the
+# writer is not closed. batch_size is larger than the number of rows and the 
checkpoint interval is long,
+# so the only thing that can flush the buffered row to Prometheus is the 
engine timer flush.
+env {
+  parallelism = 1
+  job.mode = "STREAMING"
+  checkpoint.interval = 300000
+  sink.flush.interval = 2000
+}
+
+source {
+  FakeSource {
+    schema = {
+      fields {
+        c_map = "map<string, string>"
+        c_double = double
+        c_timestamp = timestamp
+      }
+    }
+    plugin_output = "fake"
+    rows = [
+      {
+        kind = INSERT
+        fields = [{"__name__" : "timer_flush_metric"}, 2.34, CURRENT_TIMESTAMP]
+      }
+    ]
+  }
+}
+
+sink {
+  Prometheus {
+    plugin_input = "fake"
+    url = "http://prometheus-host:9090/api/v1/write";
+    key_label = "c_map"
+    key_value = "c_double"
+    key_timestamp = "c_timestamp"
+    batch_size = 100
+  }
+}

Reply via email to