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

CritasWang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb-extras.git


The following commit(s) were added to refs/heads/master by this push:
     new 40d0889  Add the TC-1 ingestion-throughput benchmark for the Table 
Mode write path (#120)
40d0889 is described below

commit 40d0889f5ee74ac1fb2c324d3f5be5a86e547637
Author: ZIHAN DAI <[email protected]>
AuthorDate: Mon Aug 10 11:57:10 2026 +1000

    Add the TC-1 ingestion-throughput benchmark for the Table Mode write path 
(#120)
    
    Adds a JUnit-driven smoke benchmark that drives the real save path -- 
bounded
    queue, single flush worker, multi-row Tablet insert -- against the same
    apache/iotdb:2.0.8-standalone Testcontainer the functional integration tests
    use, and reports records/sec, error rate and writer statistics.
    
    The test asserts a deliberately conservative floor of 1,000 rows/sec plus
    strict correctness: zero failures, zero rejects, flushed equal to rows 
written,
    and sampled rows readable back from IoTDB. Its job is to catch a gross
    throughput regression and to prove real end-to-end ingestion, not to 
certify a
    headline number -- 50 saver threads drive only 50 distinct devices, 
entity_id
    is a TAG column so device cardinality materially changes the workload, and
    30,000 rows fit entirely in the 50,000-row queue, so the run exercises no
    back-pressure. It carries the benchmark tag and is not part of the default
    test run.
    
    docs/benchmarks records the case definition, the two profiles, and the 
measured
    result with its provenance: a shared developer laptop under Docker Desktop
    rather than a dedicated benchmark host. Two back-to-back runs on that host
    measured 54,292 and 61,936 rows/sec with a zero error rate; the report 
anchors
    on the conservative figure and states explicitly that this does not 
constitute
    a pass of the sustained full-profile target, which is defined for 1,000 
devices
    on dedicated hardware and remains deferred.
    
    docker-compose.bench.yml provides the equivalent stack for running the case
    outside JUnit.
    
    Signed-off-by: Zihan Dai <[email protected]>
---
 iotdb-thingsboard-table/docker-compose.bench.yml   |  66 ++++
 iotdb-thingsboard-table/docs/benchmarks/README.md  | 170 ++++++++
 iotdb-thingsboard-table/docs/benchmarks/report.md  |  91 +++++
 .../table/IoTDBTableIngestionBenchmarkIT.java      | 431 +++++++++++++++++++++
 4 files changed, 758 insertions(+)

diff --git a/iotdb-thingsboard-table/docker-compose.bench.yml 
b/iotdb-thingsboard-table/docker-compose.bench.yml
new file mode 100644
index 0000000..016d7c1
--- /dev/null
+++ b/iotdb-thingsboard-table/docker-compose.bench.yml
@@ -0,0 +1,66 @@
+#
+# 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.
+#
+
+# TC-1 ingestion-throughput SMOKE-profile bench stack.
+#
+# This is the smoke profile: a single, clean-volume IoTDB 2.0.8 node for a 
fast,
+# reproducible local benchmark run. It mirrors docker-compose.test.yml but 
ships
+# only the IoTDB service on a dedicated fresh volume so each run starts from an
+# empty store. The benchmark IT (IoTDBTableIngestionBenchmarkIT) provisions its
+# own throwaway Testcontainer, so this compose file is only for an out-of-band
+# manual smoke run against a standalone node.
+#
+# The FULL multi-backend profile (Cassandra / PostgreSQL / TimescaleDB on a
+# dedicated host, contributor-run) is later-scope and is not defined here.
+#
+# Do not hardcode passwords or local hostnames; pass IOTDB_USERNAME /
+# IOTDB_PASSWORD via the environment.
+
+services:
+  iotdb:
+    image: apache/iotdb:2.0.8-standalone
+    container_name: iotdb-table-bench
+    environment:
+      IOTDB_USERNAME: ${IOTDB_USERNAME:?set IOTDB_USERNAME}
+      IOTDB_PASSWORD: ${IOTDB_PASSWORD:?set IOTDB_PASSWORD}
+      # Bind the client RPC service to all interfaces so the mapped host port 
is
+      # reachable from the benchmark client (default 127.0.0.1 only listens on 
the
+      # container loopback).
+      dn_rpc_address: 0.0.0.0
+    ports:
+      - "${IOTDB_RPC_PORT:-6667}:6667"
+    volumes:
+      # Fresh, dedicated volume so every smoke run starts from an empty store.
+      # Reset between runs with: docker compose -f docker-compose.bench.yml 
down -v
+      - iotdb-bench-data:/iotdb/data
+    networks:
+      - tb-iotdb-bench
+    healthcheck:
+      test: ["CMD-SHELL", "bash -ec ': >/dev/tcp/127.0.0.1/6667'"]
+      interval: 10s
+      timeout: 5s
+      retries: 12
+      start_period: 30s
+
+networks:
+  tb-iotdb-bench:
+    driver: bridge
+
+volumes:
+  iotdb-bench-data:
diff --git a/iotdb-thingsboard-table/docs/benchmarks/README.md 
b/iotdb-thingsboard-table/docs/benchmarks/README.md
new file mode 100644
index 0000000..898c0a8
--- /dev/null
+++ b/iotdb-thingsboard-table/docs/benchmarks/README.md
@@ -0,0 +1,170 @@
+<!--
+
+    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.
+
+-->
+
+# IoTDB ThingsBoard Table — Benchmarks
+
+This directory documents the performance test cases for the IoTDB Table Mode
+ThingsBoard storage backend. It currently covers **TC-1 (ingestion
+throughput)**, the write-throughput case for this backend.
+
+## Two profiles
+
+Each test case is defined with two profiles:
+
+- **Smoke profile** — small, fast, reproducible on a laptop in well under ten
+  minutes; single tenant; not wired into required CI, though it is CI-eligible
+  wherever Docker is available. Its purpose is to exercise the real write path
+  end to end and guard against gross regressions. **This is what is implemented
+  today** (`IoTDBTableIngestionBenchmarkIT`).
+- **Full profile** — a dedicated host, contributor-run, multi-backend
+  comparison (Cassandra / PostgreSQL / TimescaleDB). It is *not* in CI and is
+  **later-scope**; its report lives in [`report.md`](report.md) (placeholder
+  until a fresh run is recorded).
+
+## TC-1 — Ingestion throughput
+
+Design intent: 1,000 devices writing simultaneously via 50 concurrent threads 
in
+500-entry batches; measure records/sec and error rate. The calendar target is
+**> 10,000 writes/sec**.
+
+That **> 10K writes/sec figure is the full-profile headline on a dedicated
+host.** A cold single-node container on a laptop or CI runner will not reach 
it,
+so the smoke profile does not assert it.
+
+### What the smoke benchmark does
+
+`IoTDBTableIngestionBenchmarkIT` drives the **real** save path — the same code
+ThingsBoard uses in production:
+
+```
+dao.save(tenant, entity, tsKvEntry, ttl)
+  -> writer.enqueue(...)           bounded ArrayBlockingQueue (capacity 50,000)
+  -> single flush worker           batches up to 500 rows, maxLingerMs 20
+  -> Tablet insert                 multi-row table-session insert
+  -> real IoTDB 2.0.8              apache/iotdb:2.0.8-standalone Testcontainer
+```
+
+It runs `SAVER_THREADS = 50` concurrent threads, each writing
+`ROWS_PER_THREAD = 600` rows (30,000 rows total) with a distinct
+`(entity, key, timestamp)` per write so nothing is deduplicated away. The
+production save defaults are used unchanged (batchSize 500, queueCapacity
+50,000, maxLingerMs 20, flushThreads 1, sessionPoolSize 8); only the retry
+backoff is shortened so a transient cold-start blip does not stretch the
+measured window. The total row count is kept below the queue capacity so the
+run is free of back-pressure rejects without changing the real defaults.
+
+### What it measures and asserts
+
+Measured and logged:
+
+- **records/sec** — `totalRows / wall-clock seconds`, timed from the first
+  `save()` to all save futures completing.
+- **error rate** — `failedFutures / totalRows`.
+- **writer stats** — `dao.stats()`: `enqueued`, `flushed`, `flushFailures`,
+  `retries`, `rejectsFull`, `rejectsShutdown`, `queueDepth`.
+- **persisted-sample count** — a handful of rows are read back from IoTDB to
+  prove real ingestion, not just future completion.
+
+Asserted:
+
+- error rate `== 0` and zero failed save futures;
+- `flushFailures == 0`, `rejectsFull == 0`, `rejectsShutdown == 0`;
+- `flushed == totalRows` (every distinct row reached IoTDB);
+- the sampled rows are readable back from IoTDB;
+- throughput `>=` a **conservative smoke floor of 1,000 rows/sec**.
+
+#### Why the floor is 1,000 rows/sec, not 10,000
+
+The smoke floor only guards against gross regressions and proves correctness on
+a cold, shared, single-node container. It is intentionally an order of 
magnitude
+below the full-profile headline so the test is not flaky on laptops or CI. 
Raise
+it only alongside a measured full-profile report — never to chase the headline
+number on CI.
+
+### How to run it
+
+The benchmark is named `*IT.java`, so the unit `mvn test` run never executes 
it;
+it only test-compiles there. Integration tests in this module run via the
+**`iotdb-table-it` profile** (maven-failsafe-plugin), which requires Docker.
+
+The benchmark is tagged `@Tag("benchmark")` and `@Tag("integration")`. Use the
+JUnit tag filter to select it. Run **only** the benchmark:
+
+```bash
+cd iotdb-thingsboard-table
+mvn -ntp -Piotdb-table-it verify -Dgroups=benchmark
+```
+
+`-Dgroups=benchmark` runs only the `@Tag("benchmark")` test among the failsafe
+`**/*IT.java` set, so the functional ITs are skipped and only the throughput
+benchmark runs. Use `-Dgroups=benchmark`, **not**
+`-Dtest=IoTDBTableIngestionBenchmarkIT`: a global `-Dtest=` overrides the
+include/exclude filters of the surefire executions too, which can pull a Docker
+IT into the unit `test` phase. The tag filter is applied on top of the file
+patterns instead.
+
+The benchmark also runs as part of the normal integration-test gate
+(`mvn -ntp -Piotdb-table-it verify`): it is a deliberately cheap (~10 s, 30,000
+rows) **throughput-regression guard** that asserts only a conservative floor of
+1,000 records/sec, so it stays non-flaky on shared CI while still catching a
+write-path performance regression. It needs no benchmark-specific pom wiring. 
If
+you want a purely functional gate, exclude its tag:
+
+```bash
+cd iotdb-thingsboard-table
+mvn -ntp -Piotdb-table-it verify -DexcludedGroups=benchmark
+```
+
+> **Gate note.** A bare `mvn -Piotdb-table-it verify` runs the benchmark as one
+> of the `**/*IT.java` set (it is `@Tag("benchmark")`) — this is intended, as 
it
+> guards the throughput floor. Use `-Dgroups=benchmark` to run *only* it and 
read
+> the throughput number; use `-DexcludedGroups=benchmark` to skip it.
+
+The measured records/sec and the full writer-stats report are emitted to the
+test log at INFO and to stdout, so the figure is captured even when no SLF4J
+binding is on the test classpath.
+
+If Docker is unavailable the test is skipped
+(`@Testcontainers(disabledWithoutDocker = true)`); it never fails the build for
+lack of Docker.
+
+### Smoke stack
+
+The benchmark IT manages its own throwaway `apache/iotdb:2.0.8-standalone`
+Testcontainer, so no external stack is required to run it. For a manual run
+against a standalone node instead of the throwaway container, the module's
+[`../../docker-compose.test.yml`](../../docker-compose.test.yml) brings up an
+IoTDB service (among the full ThingsBoard test stack):
+
+```bash
+IOTDB_USERNAME=<iotdb-user> IOTDB_PASSWORD=<iotdb-password> \
+  docker compose -f docker-compose.test.yml up -d iotdb
+
+# reset to an empty store between runs
+docker compose -f docker-compose.test.yml down -v
+```
+
+## Full-profile report
+
+The full-profile multi-backend report is deferred; see
+[`report.md`](report.md). The smoke-profile records/sec from a fresh run is 
also
+filled in there by whoever runs the benchmark — the numbers are never committed
+ahead of an actual run.
diff --git a/iotdb-thingsboard-table/docs/benchmarks/report.md 
b/iotdb-thingsboard-table/docs/benchmarks/report.md
new file mode 100644
index 0000000..8388cf1
--- /dev/null
+++ b/iotdb-thingsboard-table/docs/benchmarks/report.md
@@ -0,0 +1,91 @@
+<!--
+
+    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.
+
+-->
+
+# TC-1 Ingestion-Throughput Benchmark Report
+
+Methodology is documented in [`README.md`](README.md). This file records the
+measured numbers. **No figure here is committed ahead of an actual run** — 
every
+value below is filled in from a fresh local run, not copied from a previous 
one.
+
+## Smoke profile (`IoTDBTableIngestionBenchmarkIT`)
+
+Single-node `apache/iotdb:2.0.8-standalone` Testcontainer; 50 concurrent saver
+threads; 600 rows/thread (30,000 rows total); production save defaults
+(batchSize 500, queueCapacity 50,000, maxLingerMs 20, flushThreads 1,
+sessionPoolSize 8). records/sec is timed from the first `save()` to all save
+futures completing.
+
+How to reproduce:
+
+```bash
+cd iotdb-thingsboard-table
+mvn -ntp -Piotdb-table-it verify -Dgroups=benchmark
+```
+
+Results (fresh runs, 2026-08-03):
+
+| Field | Value |
+| --- | --- |
+| Date | 2026-08-03 |
+| Host | Apple Silicon laptop (macOS) via Docker Desktop / Testcontainers — a 
shared dev host, not a dedicated benchmark machine |
+| IoTDB image | `apache/iotdb:2.0.8-standalone` (Testcontainers-managed) |
+| Total rows | 30,000 |
+| Saver threads | 50 |
+| Batch size | 500 |
+| Queue capacity | 50,000 |
+| Elapsed (s) | 0.55 |
+| **records/sec** | **54,292** |
+| Error rate | 0.0000 |
+| flushed / flushFailures / rejectsFull | 30,000 / 0 / 0 |
+| retries / rejectsShutdown / queueDepth | 0 / 0 / 0 |
+
+A second back-to-back run on the same host and configuration measured **61,936 
rows/sec** in 0.48 s (error rate 0; flushed / flushFailures / rejectsFull = 
30,000 / 0 / 0). Both runs are stable; the conservative **54,292** figure 
anchors the ≈ 5.4× / ≈ 54× ratios below.
+
+> The smoke run asserts only a conservative floor of 1,000 rows/sec on a cold,
+> shared, single-node container (so it stays non-flaky as a CI regression 
guard).
+> The observed **54,292 rows/sec** is ≈5.4× the **> 10,000 writes/sec** design
+> target and ≈54× the 1,000 rows/sec smoke floor — but treat it as a
+> regression-guard / peak figure, **not** as a pass of the sustained 
full-profile
+> target. That target is defined for 1,000 devices on a dedicated host (see the
+> full profile below); this smoke run drives only 50 distinct devices (one per
+> saver thread), and because `entity_id` is a TAG column the device cardinality
+> materially changes the write workload. The 30,000 rows also fit entirely in 
the
+> 50,000-row queue (so `rejectsFull=0`, no back-pressure) and are drained by a
+> single flush worker. The sustained 1,000-device > 10K target therefore 
remains
+> the deferred full profile, not something this smoke run validates.
+
+## Full profile (deferred / later-scope)
+
+> **Status: deferred / later-scope.**
+
+This section will hold the **full-profile** TC-1 ingestion-throughput results:
+1,000 devices, 50 concurrent threads, 500-entry batches, run on a dedicated
+host, with the **> 10,000 writes/sec** target and a multi-backend comparison
+(Cassandra / PostgreSQL / TimescaleDB).
+
+The full profile is not in CI and is run by a contributor on dedicated
+hardware. To be filled in:
+
+- Hardware and IoTDB topology (single node vs. cluster).
+- Dataset: device count, keys per device, batch size, total rows.
+- Measured records/sec, error rate, and p50 / p99 batch flush latency.
+- Per-backend comparison table.
+- Tuning notes (session pool size, flush threads, queue capacity, linger).
diff --git 
a/iotdb-thingsboard-table/src/test/java/org/apache/iotdb/extras/thingsboard/table/IoTDBTableIngestionBenchmarkIT.java
 
b/iotdb-thingsboard-table/src/test/java/org/apache/iotdb/extras/thingsboard/table/IoTDBTableIngestionBenchmarkIT.java
new file mode 100644
index 0000000..708d6b4
--- /dev/null
+++ 
b/iotdb-thingsboard-table/src/test/java/org/apache/iotdb/extras/thingsboard/table/IoTDBTableIngestionBenchmarkIT.java
@@ -0,0 +1,431 @@
+/*
+ * 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.iotdb.extras.thingsboard.table;
+
+import org.apache.iotdb.isession.ITableSession;
+import org.apache.iotdb.isession.SessionDataSet;
+import org.apache.iotdb.isession.pool.ITableSessionPool;
+import org.apache.iotdb.session.pool.TableSessionPoolBuilder;
+
+import com.google.common.util.concurrent.ListenableFuture;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.testcontainers.containers.GenericContainer;
+import org.testcontainers.containers.wait.strategy.Wait;
+import org.testcontainers.junit.jupiter.Container;
+import org.testcontainers.junit.jupiter.Testcontainers;
+import org.testcontainers.utility.DockerImageName;
+import org.thingsboard.server.common.data.EntityType;
+import org.thingsboard.server.common.data.id.EntityId;
+import org.thingsboard.server.common.data.id.TenantId;
+import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
+import org.thingsboard.server.common.data.kv.LongDataEntry;
+
+import java.io.InputStream;
+import java.nio.charset.StandardCharsets;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.UUID;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * TC-1 ingestion-throughput benchmark for the IoTDB Table Mode timeseries 
write path (smoke
+ * profile).
+ *
+ * <p>TC-1 is defined with two profiles. This class is the <b>smoke 
profile</b>: a local, fast,
+ * JUnit-driven run that exercises the real {@link 
IoTDBTableTimeseriesDao#save} path (bounded queue
+ * &rarr; single flush worker &rarr; multi-row {@code Tablet} insert &rarr; 
real IoTDB) against the
+ * same {@code apache/iotdb:2.0.8-standalone} Testcontainer the functional ITs 
use, then reports
+ * records/sec, error rate, and writer stats.
+ *
+ * <p><b>The &gt;10K writes/sec headline target is the FULL profile number on 
a dedicated host.</b>
+ * A cold single-node Testcontainer on a laptop/CI runner will not reach it, 
so this smoke test only
+ * asserts a deliberately conservative throughput floor plus strict 
correctness (zero failures, zero
+ * rejects, flushed == rows written, sample rows persisted). Its job is to 
guard against gross
+ * throughput regressions and prove real end-to-end ingestion, not to certify 
the headline number.
+ * The full multi-backend comparison (Cassandra / PostgreSQL / TimescaleDB) is 
later-scope and is
+ * not built here.
+ *
+ * @since 2.0.4-SNAPSHOT
+ */
+@Tag("benchmark")
+@Tag("integration")
+@Testcontainers(disabledWithoutDocker = true)
+class IoTDBTableIngestionBenchmarkIT {
+  private static final Logger LOG = 
LoggerFactory.getLogger(IoTDBTableIngestionBenchmarkIT.class);
+
+  // Smoke sizing: concurrency mirrors the TC-1 design (50 concurrent threads) 
but the total row
+  // count is kept modest so the run finishes in well under the 
integration-test budget on a cold
+  // container. ROWS_PER_THREAD is chosen so SAVER_THREADS * ROWS_PER_THREAD 
stays comfortably
+  // below the production save queue capacity (50_000), which keeps the run 
free of back-pressure
+  // rejects without changing the real defaults.
+  private static final int SAVER_THREADS = 50;
+  private static final int ROWS_PER_THREAD = 600;
+  private static final int TOTAL_ROWS = SAVER_THREADS * ROWS_PER_THREAD; // 
30_000
+
+  // Conservative smoke floor. The >10K rows/sec design-doc target is the 
FULL-profile headline on a
+  // dedicated host; a cold single-node Testcontainer cannot be held to it 
without flakiness, so we
+  // only assert that the real save() path sustains at least this floor 
end-to-end. Raise this only
+  // alongside a measured full-profile report (docs/benchmarks/report.md), 
never to chase the
+  // headline on CI.
+  private static final double SMOKE_THROUGHPUT_FLOOR_ROWS_PER_SEC = 1_000.0D;
+
+  private static final int FUTURE_TIMEOUT_SECONDS = 60;
+  // One global ceiling for awaiting the WHOLE set of save futures. A systemic 
writer stall must
+  // fail
+  // the smoke benchmark within this bound instead of applying a per-future 
timeout to each of
+  // TOTAL_ROWS futures in turn (which would let a hang run for hours before 
CI kills it).
+  private static final int AWAIT_ALL_TIMEOUT_SECONDS = 120;
+  private static final int VERIFY_SAMPLE_KEYS = 5;
+  private static final Duration IOTDB_STARTUP_TIMEOUT = Duration.ofMinutes(3);
+  private static final Duration IOTDB_READY_TIMEOUT = Duration.ofSeconds(60);
+  private static final Duration IOTDB_READY_POLL_INTERVAL = 
Duration.ofMillis(500);
+
+  @Container
+  static final GenericContainer<?> IOTDB =
+      new 
GenericContainer<>(DockerImageName.parse("apache/iotdb:2.0.8-standalone"))
+          .withExposedPorts(6667)
+          // IoTDB binds its client RPC service to dn_rpc_address (default 
127.0.0.1), so it would
+          // only listen on the container loopback and reject the 
Testcontainers port-mapped session
+          // handshake ("Fail to reconnect"). Bind to all interfaces so the 
mapped host port works.
+          .withEnv("dn_rpc_address", "0.0.0.0")
+          
.waitingFor(Wait.forListeningPort().withStartupTimeout(IOTDB_STARTUP_TIMEOUT));
+
+  @Test
+  void tc1_ingestionThroughput_smokeProfile() throws Exception {
+    BenchmarkScope scope = scope();
+    bootstrapSchema(scope.database());
+    try (ITableSessionPool pool = newPool(scope.database())) {
+      IoTDBTableConfig config = benchmarkConfig();
+      IoTDBTableTimeseriesWriter writer = new IoTDBTableTimeseriesWriter(pool, 
config);
+      IoTDBTableTimeseriesDao dao = new IoTDBTableTimeseriesDao(pool, writer, 
config);
+
+      ExecutorService savers = Executors.newFixedThreadPool(SAVER_THREADS, 
saverThreadFactory());
+      List<ListenableFuture<Integer>> futures = new ArrayList<>(TOTAL_ROWS);
+      AtomicInteger failedSubmits = new AtomicInteger();
+      CountDownLatch ready = new CountDownLatch(SAVER_THREADS);
+      CountDownLatch start = new CountDownLatch(1);
+      CountDownLatch done = new CountDownLatch(SAVER_THREADS);
+
+      try {
+        // Each thread owns a disjoint entity + timestamp band so every 
(entity, key, ts) tuple is
+        // unique across the whole run; nothing is deduplicated away and 
flushed must equal
+        // TOTAL_ROWS.
+        for (int threadIndex = 0; threadIndex < SAVER_THREADS; threadIndex++) {
+          int t = threadIndex;
+          savers.execute(
+              () -> {
+                EntityId entity = entityForThread(t);
+                List<ListenableFuture<Integer>> local = new 
ArrayList<>(ROWS_PER_THREAD);
+                ready.countDown();
+                try {
+                  start.await();
+                  for (int r = 0; r < ROWS_PER_THREAD; r++) {
+                    long ts = ((long) t * ROWS_PER_THREAD) + r + 1L;
+                    BasicTsKvEntry entry =
+                        new BasicTsKvEntry(ts, new LongDataEntry("metric", 
(long) r));
+                    local.add(dao.save(scope.tenantId(), entity, entry, 0));
+                  }
+                } catch (InterruptedException e) {
+                  Thread.currentThread().interrupt();
+                } catch (RuntimeException e) {
+                  failedSubmits.incrementAndGet();
+                  LOG.warn("TC-1 saver thread {} failed to submit", t, e);
+                } finally {
+                  synchronized (futures) {
+                    futures.addAll(local);
+                  }
+                  done.countDown();
+                }
+              });
+        }
+
+        assertTrue(ready.await(30, TimeUnit.SECONDS), "saver threads did not 
become ready in time");
+        long startNanos = System.nanoTime();
+        start.countDown();
+        assertTrue(
+            done.await(FUTURE_TIMEOUT_SECONDS, TimeUnit.SECONDS),
+            "saver threads did not finish submitting in time");
+
+        assertEquals(0, failedSubmits.get(), "save() submission must not 
throw");
+        assertEquals(TOTAL_ROWS, futures.size(), "every row must produce a 
save future");
+
+        long failedFutures = awaitAll(futures);
+        long elapsedNanos = System.nanoTime() - startNanos;
+        double elapsedSeconds = elapsedNanos / 1_000_000_000.0D;
+        double recordsPerSec = TOTAL_ROWS / elapsedSeconds;
+        double errorRate = (double) failedFutures / TOTAL_ROWS;
+
+        IoTDBTableTimeseriesWriterStats stats = dao.stats();
+        long persistedSample = countPersistedSample(pool, scope);
+
+        String report =
+            String.format(
+                "TC-1 ingestion smoke profile: rows=%d threads=%d batchSize=%d 
queueCapacity=%d "
+                    + "elapsed=%.2fs recordsPerSec=%.0f errorRate=%.4f | 
stats[enqueued=%d "
+                    + "flushed=%d flushFailures=%d retries=%d rejectsFull=%d 
rejectsShutdown=%d "
+                    + "queueDepth=%d] | persistedSampleRows=%d (across %d 
sampled threads) "
+                    + "| floor=%.0frows/sec (NOTE: >10K rows/sec is the 
FULL-profile headline on a "
+                    + "dedicated host; this smoke run only guards 
regressions)",
+                TOTAL_ROWS,
+                SAVER_THREADS,
+                config.getTs().getSave().getBatchSize(),
+                config.getTs().getSave().getQueueCapacity(),
+                elapsedSeconds,
+                recordsPerSec,
+                errorRate,
+                stats.enqueued(),
+                stats.flushed(),
+                stats.flushFailures(),
+                stats.retries(),
+                stats.rejectsFull(),
+                stats.rejectsShutdown(),
+                stats.queueDepth(),
+                persistedSample,
+                VERIFY_SAMPLE_KEYS,
+                SMOKE_THROUGHPUT_FLOOR_ROWS_PER_SEC);
+        LOG.info(report);
+        // Also emit to stdout so the measured records/sec is captured in the 
surefire/failsafe
+        // console output even when no SLF4J binding is on the test classpath 
(NOP logger).
+        System.out.println(report);
+
+        // Correctness: the real save path must complete every row with no 
failures or rejects.
+        assertEquals(0L, failedFutures, "TC-1 smoke profile must complete with 
zero failed saves");
+        assertEquals(0.0D, errorRate, "TC-1 smoke profile error rate must be 
zero");
+        assertEquals(0L, stats.flushFailures(), "writer flushFailures must be 
zero");
+        assertEquals(
+            0L, stats.rejectsFull(), "writer rejectsFull must be zero (queue 
not saturated)");
+        assertEquals(0L, stats.rejectsShutdown(), "writer rejectsShutdown must 
be zero");
+        assertEquals(
+            TOTAL_ROWS,
+            stats.flushed(),
+            "every distinct row must be flushed (nothing deduplicated)");
+
+        // Proof of real ingestion: a sample of rows must be readable back 
from IoTDB.
+        assertEquals(
+            VERIFY_SAMPLE_KEYS,
+            persistedSample,
+            "sampled rows must be persisted and readable from IoTDB");
+
+        // Conservative regression floor, not the >10K rows/sec full-profile 
headline.
+        assertTrue(
+            recordsPerSec >= SMOKE_THROUGHPUT_FLOOR_ROWS_PER_SEC,
+            () ->
+                "TC-1 smoke throughput "
+                    + String.format("%.0f", recordsPerSec)
+                    + " rows/sec fell below the conservative smoke floor "
+                    + String.format("%.0f", 
SMOKE_THROUGHPUT_FLOOR_ROWS_PER_SEC)
+                    + " rows/sec (full-profile target is >10K on a dedicated 
host)");
+      } finally {
+        savers.shutdownNow();
+        try {
+          // Best-effort: on an early assertion failure, let interrupted saver 
threads unwind before
+          // we tear down the DAO/writer they may still be calling into.
+          savers.awaitTermination(10, TimeUnit.SECONDS);
+        } catch (InterruptedException e) {
+          Thread.currentThread().interrupt();
+        }
+        dao.destroy();
+        writer.destroy();
+      }
+    }
+  }
+
+  private long awaitAll(List<ListenableFuture<Integer>> futures) throws 
InterruptedException {
+    // Single shared deadline for the whole set: once it passes, each 
remaining future is polled
+    // with
+    // a zero (non-blocking) budget, so a stall is detected fast and the total 
wait is bounded by
+    // AWAIT_ALL_TIMEOUT_SECONDS regardless of how many futures are 
outstanding.
+    long deadlineNanos = System.nanoTime() + 
TimeUnit.SECONDS.toNanos(AWAIT_ALL_TIMEOUT_SECONDS);
+    long failed = 0L;
+    Throwable firstFailure = null;
+    for (ListenableFuture<Integer> future : futures) {
+      long remainingNanos = Math.max(0L, deadlineNanos - System.nanoTime());
+      try {
+        future.get(remainingNanos, TimeUnit.NANOSECONDS);
+      } catch (java.util.concurrent.ExecutionException | 
java.util.concurrent.TimeoutException e) {
+        failed++;
+        if (firstFailure == null) {
+          firstFailure = e;
+        }
+      }
+    }
+    if (failed > 0L) {
+      LOG.warn(
+          "TC-1: {} of {} save futures did not complete within {}s (first 
failure shown)",
+          failed,
+          futures.size(),
+          AWAIT_ALL_TIMEOUT_SECONDS,
+          firstFailure);
+    }
+    return failed;
+  }
+
+  /**
+   * Reads back the first row written by the first {@link #VERIFY_SAMPLE_KEYS} 
threads to prove the
+   * benchmark persisted real rows rather than merely completing futures.
+   */
+  private long countPersistedSample(ITableSessionPool pool, BenchmarkScope 
scope) throws Exception {
+    long found = 0L;
+    for (int t = 0; t < VERIFY_SAMPLE_KEYS; t++) {
+      EntityId entity = entityForThread(t);
+      long ts = ((long) t * ROWS_PER_THREAD) + 1L;
+      String sql =
+          "SELECT long_v FROM telemetry WHERE tenant_id='"
+              + scope.tenantId().getId()
+              + "' AND entity_type='DEVICE' AND entity_id='"
+              + entity.getId()
+              + "' AND key='metric' AND time="
+              + ts;
+      try (ITableSession session = pool.getSession();
+          SessionDataSet dataSet = session.executeQueryStatement(sql)) {
+        SessionDataSet.DataIterator row = dataSet.iterator();
+        if (row.next() && !row.isNull("long_v")) {
+          found++;
+        }
+      }
+    }
+    return found;
+  }
+
+  private EntityId entityForThread(int threadIndex) {
+    // Deterministic per-thread device UUID so each thread targets a distinct 
entity.
+    UUID id = new UUID(0xBE0000000000L, 0x1000L + threadIndex);
+    return new BenchmarkEntityId(id);
+  }
+
+  private IoTDBTableConfig benchmarkConfig() {
+    // Production save-path defaults: batchSize=500, queueCapacity=50000,
+    // maxLingerMs=20, flushThreads=1, sessionPoolSize=8. Only the retry 
backoff is shortened so a
+    // transient cold-start blip does not stretch the measured window; the 
throughput-relevant
+    // knobs are left at their real defaults so the smoke run exercises the 
real configuration.
+    IoTDBTableConfig config = new IoTDBTableConfig();
+    config.getTs().getSave().setRetryInitialBackoffMs(1L);
+    config.getTs().getSave().setRetryMaxBackoffMs(1L);
+    config.getTs().getRead().setThreads(1);
+    return config;
+  }
+
+  private ITableSessionPool newPool(String database) {
+    TableSessionPoolBuilder builder =
+        new TableSessionPoolBuilder()
+            .nodeUrls(List.of("127.0.0.1:" + IOTDB.getMappedPort(6667)))
+            .user("root")
+            .password("root")
+            .maxSize(8);
+    if (database != null) {
+      builder.database(database);
+    }
+    return builder.build();
+  }
+
+  private void bootstrapSchema(String database) throws Exception {
+    awaitIoTDBReady(database);
+
+    String schema;
+    try (InputStream stream =
+        IoTDBTableIngestionBenchmarkIT.class
+            .getClassLoader()
+            .getResourceAsStream("schema-iotdb-table.sql")) {
+      schema = new String(stream.readAllBytes(), StandardCharsets.UTF_8);
+    }
+    schema =
+        schema
+            .replace(
+                "CREATE DATABASE IF NOT EXISTS thingsboard;",
+                "CREATE DATABASE IF NOT EXISTS " + database + ";")
+            .replace("USE thingsboard;", "USE " + database + ";");
+    schema = schema.replaceAll("(?s)/\\*.*?\\*/", "").replaceAll("(?m)--.*$", 
"");
+    try (ITableSessionPool bootstrapPool = newPool(null);
+        ITableSession session = bootstrapPool.getSession()) {
+      for (String statement : schema.split(";")) {
+        String trimmed = statement.trim();
+        if (!trimmed.isEmpty()) {
+          session.executeNonQueryStatement(trimmed);
+        }
+      }
+    }
+  }
+
+  private void awaitIoTDBReady(String database) throws Exception {
+    long deadlineNanos = System.nanoTime() + IOTDB_READY_TIMEOUT.toNanos();
+    Exception lastFailure = null;
+    while (System.nanoTime() < deadlineNanos) {
+      try (ITableSessionPool bootstrapPool = newPool(null);
+          ITableSession session = bootstrapPool.getSession()) {
+        session.executeNonQueryStatement("CREATE DATABASE IF NOT EXISTS " + 
database);
+        return;
+      } catch (Exception e) {
+        lastFailure = e;
+        long remainingMillis = TimeUnit.NANOSECONDS.toMillis(deadlineNanos - 
System.nanoTime());
+        if (remainingMillis <= 0) {
+          break;
+        }
+        Thread.sleep(Math.min(IOTDB_READY_POLL_INTERVAL.toMillis(), 
remainingMillis));
+      }
+    }
+    throw new IllegalStateException(
+        "IoTDB did not accept table-session statements within " + 
IOTDB_READY_TIMEOUT, lastFailure);
+  }
+
+  private BenchmarkScope scope() {
+    return new BenchmarkScope(
+        uniqueDatabase(), new 
TenantId(UUID.fromString("55555555-5555-5555-5555-555555555501")));
+  }
+
+  private String uniqueDatabase() {
+    // IoTDB caps database names at 64 chars; keep the prefix short and append 
a trimmed UUID.
+    String shortUuid = UUID.randomUUID().toString().replace("-", 
"").substring(0, 16);
+    return "tb_bench_tc1_" + shortUuid;
+  }
+
+  private static java.util.concurrent.ThreadFactory saverThreadFactory() {
+    AtomicLong sequence = new AtomicLong();
+    return runnable -> {
+      Thread thread = new Thread(runnable, "tc1-benchmark-saver-" + 
sequence.incrementAndGet());
+      thread.setDaemon(true);
+      return thread;
+    };
+  }
+
+  private record BenchmarkScope(String database, TenantId tenantId) {}
+
+  private record BenchmarkEntityId(UUID id) implements EntityId {
+    @Override
+    public UUID getId() {
+      return id;
+    }
+
+    @Override
+    public EntityType getEntityType() {
+      return EntityType.DEVICE;
+    }
+  }
+}

Reply via email to