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]

Reply via email to