albgen opened a new issue, #12022: URL: https://github.com/apache/seatunnel/issues/12022
### Search before asking - [x] I had searched in the [feature](https://github.com/apache/seatunnel/issues?q=is%3Aissue+label%3A%22Feature%22) and found no similar feature requirement. ### Description ### Description Zeta persists every `engine*` IMap through `FileMapStoreFactory`. That set includes the two metrics maps, `engine_runningJobMetrics` and `engine_finishedJobMetrics`, and they dominate the write volume by a wide margin: the value is the whole `HashMap<TaskLocation, SeaTunnelMetricsContext>` of a job, re-put on every metrics flush, and the file store appends the full serialized value each time. The store is append-only and never compacted (#10329), and with `initial-mode: EAGER` all of it is replayed at every start. Measured on a single Zeta node running 40 streaming SQL Server-CDC → Doris jobs (`parallelism = 1`, `checkpoint.interval = 5000`), after ~8 days of uptime: | map | store size | growth | |---|---|---| | `engine_runningJobMetrics` | **13 GB** (one `wal.txt` segment reached 11.95 GB) | ~65 MB/h ≈ 1.5 GB/day | | `engine_checkpoint-id-map` | 492 MB | ~60 MB/day | | all other `engine*` maps combined | < 6 MB | negligible | The live content of the metrics map is ~40 entries. The 13 GB is history that only exists to be replayed once, at startup, into the heap — where in our case it caused an `OutOfMemoryError` inside `MapProxySupport.initializeMapStoreLoad` during `CoordinatorService.restoreAllRunningJobFromMasterNodeSwitch`, leaving a node that never opens its REST port while the container still reports `Up` (full report in #10329). **Nothing restores from these two maps.** They are telemetry: live values are served from `/job-info/<jobId>` and the Prometheus endpoint, and after a restore the counters start from zero anyway. Persisting them buys nothing and is the single largest contributor to the growth described in #10329. ### Proposal Exclude the metrics maps from the persisted set by default — either in the shipped `config/hazelcast.yaml` (an exact map name beats the `engine*` wildcard, the same mechanism already used for `engine_checkpoint_monitor`), or by not registering a map-store for those maps in the server, with an opt-in flag for anyone who wants metrics history to survive a restart. This is complementary to #10329 / #10399, not a replacement — compaction is still needed for the maps that legitimately have to be persisted. But it is a config-level change that removes ~95% of the growth immediately and could ship well before LSM compaction lands. **Result in production:** with the two overrides applied, the store went from 13 GB to ~4 MB, restarts restore all 40 jobs from their checkpoints normally, and `/job-info` metrics are unaffected. --- ### SeaTunnel Version 2.3.13 --- ### SeaTunnel Config `config/hazelcast.yaml` — the map-store section (the subject of this issue). The last two entries are the proposed fix, applied locally: ```yaml hazelcast: cluster-name: seatunnel network: join: tcp-ip: enabled: true member-list: - localhost port: auto-increment: false port: 5801 map: engine*: map-store: enabled: true initial-mode: EAGER factory-class-name: org.apache.seatunnel.engine.server.persistence.FileMapStoreFactory properties: type: hdfs namespace: /opt/seatunnel/state/imap clusterName: seatunnel-cluster storage.type: hdfs fs.defaultFS: file:/// # proposed default (applied locally): metrics maps are telemetry, nothing restores from them engine_runningJobMetrics: map-store: enabled: false engine_finishedJobMetrics: map-store: enabled: false ``` A representative job (40 of these run concurrently, one per source table): ```hocon env { job.mode = "STREAMING" job.name = "cdc_customer" parallelism = 1 checkpoint.interval = 5000 checkpoint.timeout = 600000 } source { SqlServer-CDC { plugin_output = "src" base-url = "jdbc:sqlserver://<host>:1433;databaseName=<db>" username = "<user>" password = "<password>" database-names = ["<db>"] table-names = ["<db>.dbo.Customer"] startup.mode = "initial" exactly_once = false } } sink { Doris { plugin_input = "src" fenodes = "<fe>:8030" database = "raw" table = "customer" sink.enable-2pc = true sink.label-prefix = "cdc_customer_<runid>" } } ``` --- ### Running Command ```shell # cluster: single node, default role master_and_worker (container entrypoint) bin/seatunnel-cluster.sh # jobs: 40 submissions over REST, one per table curl -sS -X POST -H 'Content-Type: text/plain; charset=utf-8' \ --data-binary @cdc_customer.conf \ 'http://127.0.0.1:8080/submit-job?jobName=cdc_customer&format=hocon' ``` Then simply leave it running. `du -sh state/imap/seatunnel-cluster/*` shows the growth: ```shell 13G state/imap/seatunnel-cluster/engine_runningJobMetrics 492M state/imap/seatunnel-cluster/engine_checkpoint-id-map 5.0M state/imap/seatunnel-cluster/engine_runningJobInfo 604K state/imap/seatunnel-cluster/engine_runningJobState 556K state/imap/seatunnel-cluster/engine_stateTimestamps 580K state/imap/seatunnel-cluster/engine_finishedJobMetrics 140K state/imap/seatunnel-cluster/engine_ownedSlotProfilesIMap ``` --- ### Error Exception Growth alone produces no exception — it is silent until the next restart, when the EAGER load of the accumulated store runs out of heap (`-Xmx5g` here): ```log java.lang.OutOfMemoryError: Java heap space Dumping heap to /tmp/seatunnel/dump/zeta-server/java_pid13.hprof ... /opt/seatunnel/bin/seatunnel-cluster.sh: line 201: 13 Killed java ${JAVA_OPTS} -cp ${CLASS_PATH} ${APP_MAIN} ${args} at org.apache.seatunnel.engine.server.CoordinatorService.restoreAllRunningJobFromMasterNodeSwitch(CoordinatorService.java:499) at org.apache.seatunnel.engine.server.CoordinatorService.restoreJobFromMasterActiveSwitch(CoordinatorService.java:528) at com.hazelcast.map.impl.proxy.MapProxySupport.initializeMapStoreLoad(MapProxySupport.java:330) at com.hazelcast.instance.impl.OutOfMemoryErrorDispatcher.onOutOfMemory(OutOfMemoryErrorDispatcher.java:187) com.hazelcast.core.HazelcastInstanceNotActiveException: Hazelcast instance is not active! WARN [h.m.i.r.BasicRecordStoreLoader] - [localhost]:5801 [seatunnel] [5.1] Could not load keys from map store ``` --- ### Zeta or Flink or Spark Version Zeta (SeaTunnel Engine) 2.3.13, single node, `master_and_worker`. Embedded Hazelcast IMDG 5.1. --- ### Java or Scala Version OpenJDK 1.8.0_342 (as shipped in the `apache/seatunnel:2.3.13` image), `-Xms1g -Xmx5g`. --- ### Usage Scenario _No response_ ### Related issues _No response_ ### Are you willing to submit a PR? - [x] Yes I am willing to submit a PR! ### Code of Conduct - [x] I agree to follow this project's [Code of Conduct](https://www.apache.org/foundation/policies/conduct) -- 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]
