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 51af725848 [Improve][Connector-V2] Split Bigtable source by
sampleRowKeys for parallel reads (#11876)
51af725848 is described below
commit 51af7258487c58aacec2ce8eaf30d86d807b6c34
Author: 雷炯 <[email protected]>
AuthorDate: Sun Aug 23 22:12:12 2026 +0800
[Improve][Connector-V2] Split Bigtable source by sampleRowKeys for parallel
reads (#11876)
Co-authored-by: Cursor <[email protected]>
---
.../changelog/connector-google-bigtable.md | 1 +
docs/en/connectors/source/GoogleBigtable.md | 6 +-
.../introduction/concepts/incompatible-changes.md | 7 +
.../changelog/connector-google-bigtable.md | 1 +
docs/zh/connectors/source/GoogleBigtable.md | 6 +-
.../introduction/concepts/incompatible-changes.md | 7 +
.../seatunnel/bigtable/client/BigtableClient.java | 22 ++
.../source/BigtableSourceSplitEnumerator.java | 325 ++++++++++++++++++---
.../source/BigtableSourceSplitEnumeratorTest.java | 312 +++++++++++++++++++-
.../source/BigtableSourceStateRecoveryTest.java | 87 ++++--
10 files changed, 696 insertions(+), 78 deletions(-)
diff --git a/docs/en/connectors/changelog/connector-google-bigtable.md
b/docs/en/connectors/changelog/connector-google-bigtable.md
index 3fc3a897f9..234a6ddfab 100644
--- a/docs/en/connectors/changelog/connector-google-bigtable.md
+++ b/docs/en/connectors/changelog/connector-google-bigtable.md
@@ -2,6 +2,7 @@
| Change | Commit | Version |
| --- | --- | --- |
+|[Improve][Connector-V2] Split Bigtable source by sampleRowKeys for parallel
reads|-|dev|
|[Feature][Connector-V2] Add Google Cloud Bigtable Source and Sink
connector|https://github.com/apache/seatunnel/commit/8e57c04|dev|
</details>
diff --git a/docs/en/connectors/source/GoogleBigtable.md
b/docs/en/connectors/source/GoogleBigtable.md
index f09218a8da..b0417360a6 100644
--- a/docs/en/connectors/source/GoogleBigtable.md
+++ b/docs/en/connectors/source/GoogleBigtable.md
@@ -18,12 +18,12 @@ Reads data from Google Cloud Bigtable using the native
Bigtable Data v2 Java cli
- [ ] [stream](../../introduction/concepts/connector-v2-features.md)
- [ ] [exactly-once](../../introduction/concepts/connector-v2-features.md)
- [ ] [column projection](../../introduction/concepts/connector-v2-features.md)
-- [ ] [parallelism](../../introduction/concepts/connector-v2-features.md)
+- [x] [parallelism](../../introduction/concepts/connector-v2-features.md)
- [ ] [cdc](../../introduction/concepts/connector-v2-features.md)
:::tip
-The source is bounded. It creates one split for the configured table or
row-key range, so increasing job parallelism does not split one Bigtable scan
into multiple tablet-range reads. Each scan reads every requested cell for the
configured row range and emits one SeaTunnel row per Bigtable row.
+The source is bounded. The enumerator calls Bigtable `sampleRowKeys` to cut
the table (or the configured `start_rowkey` / `end_rowkey` range) into
tablet-sized splits, then assigns them by `hash(splitId) % parallelism`. Set
`env.parallelism` (or the source parallelism) greater than 1 so multiple
readers scan different key ranges. If sampling fails, returns no keys, or
intersects to an empty range, the connector falls back to a single split so the
job can still run. Each scan reads every [...]
:::
@@ -90,7 +90,7 @@ Maximum number of cell versions to return per column
qualifier. Default `1` retu
### scan_row_limit [int]
-Maximum number of rows to return. `-1` (default) means no limit. Use this
option together with `start_rowkey` / `end_rowkey` to do paginated full-table
scans across multiple jobs.
+Maximum number of rows to return **per split**. `-1` (default) means no limit.
When the enumerator produces multiple splits, the job-level upper bound is
about `scan_row_limit × split count`, not a single table-wide cap. Use this
option together with `start_rowkey` / `end_rowkey` to do paginated full-table
scans across multiple jobs.
### common options
diff --git a/docs/en/introduction/concepts/incompatible-changes.md
b/docs/en/introduction/concepts/incompatible-changes.md
index 9b239fc1fe..3f04c6280c 100644
--- a/docs/en/introduction/concepts/incompatible-changes.md
+++ b/docs/en/introduction/concepts/incompatible-changes.md
@@ -113,6 +113,13 @@ You need to check this document before you upgrade to
related version.
- **Mixed-version directories**: Re-materialize the directory so every
file is produced by the new version, or write pre- and post-upgrade files into
separate directories and read them independently.
- **Case-sensitive consumers**: Configure the reader for case-insensitive
schema evolution where supported, or remap the column at read time.
- **Case-only sibling fields** (for example `MD5` and `md5` in the same
struct): now representable; case-insensitive downstream consumers (such as
Hive) may treat them as ambiguous — disambiguate at the source if needed.
+
+- **Breaking Change: Google Bigtable Source `scan_row_limit` is now a
per-split cap**
+ - **Affected component**: `seatunnel-connectors-v2/connector-google-bigtable`
+ - **Description**: The enumerator now partitions a table (or the configured
`start_rowkey` / `end_rowkey` range) into tablet-sized splits via
`sampleRowKeys`. `scan_row_limit` is still applied with `query.limit(...)` once
per split in the reader. Before this change the source always produced exactly
one split, so `scan_row_limit` acted as a table-wide row cap. After this change
a table with multiple tablets yields multiple splits even when `parallelism =
1` (the single reader is assign [...]
+ - **Impact**: Existing jobs that set `scan_row_limit` to bound total output
(sampling, testing, cost control, or downstream capacity) can read far more
rows after upgrade with no config change.
+ - **Migration Guide**: If you need a table-wide cap, narrow the scan with
`start_rowkey` / `end_rowkey`, or lower `scan_row_limit` so that
`scan_row_limit × expected split count` stays within the previous budget. To
keep the previous single-split behavior, the connector still falls back to one
split when sampling fails, returns no keys, or the intersection is empty — that
is not a supported way to pin the old cap. (#11876)
+
- **Breaking Change: Iceberg Connector — source table primary key is no longer
silently inherited**
- **Affected component**: `seatunnel-connectors-v2/connector-iceberg`
- **Description**: `SchemaUtils.toIcebergSchema()` previously fell back to
the CDC source
diff --git a/docs/zh/connectors/changelog/connector-google-bigtable.md
b/docs/zh/connectors/changelog/connector-google-bigtable.md
index 3fc3a897f9..234a6ddfab 100644
--- a/docs/zh/connectors/changelog/connector-google-bigtable.md
+++ b/docs/zh/connectors/changelog/connector-google-bigtable.md
@@ -2,6 +2,7 @@
| Change | Commit | Version |
| --- | --- | --- |
+|[Improve][Connector-V2] Split Bigtable source by sampleRowKeys for parallel
reads|-|dev|
|[Feature][Connector-V2] Add Google Cloud Bigtable Source and Sink
connector|https://github.com/apache/seatunnel/commit/8e57c04|dev|
</details>
diff --git a/docs/zh/connectors/source/GoogleBigtable.md
b/docs/zh/connectors/source/GoogleBigtable.md
index f0eebfe9cd..0b360ab97d 100644
--- a/docs/zh/connectors/source/GoogleBigtable.md
+++ b/docs/zh/connectors/source/GoogleBigtable.md
@@ -18,12 +18,12 @@ import ChangeLog from
'../changelog/connector-google-bigtable.md';
- [ ] [流处理](../../introduction/concepts/connector-v2-features.md)
- [ ] [精确一次](../../introduction/concepts/connector-v2-features.md)
- [ ] [列投影](../../introduction/concepts/connector-v2-features.md)
-- [ ] [并行度](../../introduction/concepts/connector-v2-features.md)
+- [x] [并行度](../../introduction/concepts/connector-v2-features.md)
- [ ] [cdc](../../introduction/concepts/connector-v2-features.md)
:::tip
-该 Source 是有界读取。当前只会为配置的表或行键范围生成一个切分,所以提高作业并行度不会把一次 Bigtable 扫描拆成多个 tablet
范围并发读取。每次扫描会读取所请求行范围内的全部 Cell,并为每个 Bigtable 行输出一条 SeaTunnel 记录。
+该 Source 是有界读取。Enumerator 会调用 Bigtable `sampleRowKeys`,把整表(或配置的 `start_rowkey`
/ `end_rowkey` 范围)按 tablet 切成多个 split,再按 `hash(splitId) % parallelism` 分配。将
`env.parallelism`(或 Source 并行度)设为大于 1,多个 Reader 会扫描不同 key
range。若采样失败、返回空,或与用户范围求交后没有任何有效区间,则回退为单个 split,作业仍可运行。每次扫描会读取所请求行范围内的全部
Cell,并为每个 Bigtable 行输出一条 SeaTunnel 记录。
:::
@@ -90,7 +90,7 @@ Cell 时间戳过滤的结束值,不包含该时间戳,单位是微秒。
### scan_row_limit [int]
-最多读取的行数。默认值 `-1` 表示不限制。把 `scan_row_limit` 与 `start_rowkey` / `end_rowkey`
配合,可以在多个作业之间分页扫描整张表。
+每个 split 最多读取的行数。默认值 `-1` 表示不限制。当 Enumerator 切出多个 split 时,作业级上限约为
`scan_row_limit × split 数`,而不是整表一条上限。把 `scan_row_limit` 与 `start_rowkey` /
`end_rowkey` 配合,可以在多个作业之间分页扫描整张表。
### common options
diff --git a/docs/zh/introduction/concepts/incompatible-changes.md
b/docs/zh/introduction/concepts/incompatible-changes.md
index 8fbb16a4f6..130eed2408 100644
--- a/docs/zh/introduction/concepts/incompatible-changes.md
+++ b/docs/zh/introduction/concepts/incompatible-changes.md
@@ -107,6 +107,13 @@
- **混合版本目录**:将目录重新物化,使所有文件都由新版本写入;或将旧版本与新版本文件分别写入不同目录,独立读取。
- **大小写敏感的下游**:在支持的情况下将 Reader 配置为大小写不敏感的 Schema Evolution,或在读取时重映射该列。
- **仅大小写不同的同名兄弟字段**(例如同一 struct 中同时存在 `MD5` 和 `md5`):现在可被表达;大小写不敏感的下游(如
Hive)可能将其视为歧义字段,如需保留请在源头消歧。
+
+- **破坏性变更:Google Bigtable Source 的 `scan_row_limit` 变为每个 split 的上限**
+ - **影响范围**:`seatunnel-connectors-v2/connector-google-bigtable`
+ - **变更说明**:Enumerator 现在通过 `sampleRowKeys` 按 tablet 边界把表(或配置的 `start_rowkey`
/ `end_rowkey` 区间)切成多个 split。Reader 仍对每个 split 调用一次 `query.limit(...)`。此前
Source 始终只产生 1 个 split,因此 `scan_row_limit` 等价于整表行数上限。升级后,只要表有多个 tablet,即使
`parallelism = 1`(唯一 reader 会拿到全部 split),作业级上限约为 `scan_row_limit × split 数`。详见
[Google Bigtable
Source](../../connectors/source/GoogleBigtable.md#scan_row_limit-int)。
+ - **影响**:依赖 `scan_row_limit`
限制总输出量的存量作业(抽样、测试、成本控制、下游容量)在升级后、配置不变的情况下,可能读出远超以前的行数。
+ - **迁移指南**:若仍需要整表级上限,请用 `start_rowkey` / `end_rowkey` 收窄扫描范围,或下调
`scan_row_limit`,使 `scan_row_limit × 预期 split 数` 不超过原预算。采样失败、无采样点或求交为空时仍会回退为单个
split,但这不是用来锁定旧语义的受支持方式。(#11876)
+
- **破坏性变更:Iceberg 连接器 — 不再自动继承源表主键**
- **影响范围**:`seatunnel-connectors-v2/connector-iceberg`
- **变更说明**:当未显式配置 `iceberg.table.primary-keys`
时,`SchemaUtils.toIcebergSchema()`
diff --git
a/seatunnel-connectors-v2/connector-google-bigtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/bigtable/client/BigtableClient.java
b/seatunnel-connectors-v2/connector-google-bigtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/bigtable/client/BigtableClient.java
index d4d21d9718..b578d72273 100644
---
a/seatunnel-connectors-v2/connector-google-bigtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/bigtable/client/BigtableClient.java
+++
b/seatunnel-connectors-v2/connector-google-bigtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/bigtable/client/BigtableClient.java
@@ -26,6 +26,7 @@ import com.google.auth.oauth2.ServiceAccountCredentials;
import com.google.cloud.bigtable.data.v2.BigtableDataClient;
import com.google.cloud.bigtable.data.v2.BigtableDataSettings;
import com.google.cloud.bigtable.data.v2.models.BulkMutation;
+import com.google.cloud.bigtable.data.v2.models.KeyOffset;
import com.google.cloud.bigtable.data.v2.models.Mutation;
import com.google.cloud.bigtable.data.v2.models.RowMutation;
import com.google.protobuf.ByteString;
@@ -140,6 +141,27 @@ public class BigtableClient implements Serializable,
AutoCloseable {
return dataClient;
}
+ /**
+ * Samples approximate tablet-boundary row keys used by the enumerator to
build parallel source
+ * splits.
+ *
+ * <p>The returned keys delimit roughly equal-size ranges; the last key is
typically empty,
+ * meaning the table end.
+ *
+ * @return sampled tablet-boundary keys
+ * @throws BigtableConnectorException if the RPC fails
+ */
+ public List<KeyOffset> sampleRowKeys() {
+ try {
+ return dataClient.sampleRowKeys(parameters.getTable());
+ } catch (Exception e) {
+ throw new BigtableConnectorException(
+ BigtableConnectorErrorCode.TABLE_QUERY_FAILED,
+ "Failed to sample row keys for table " +
parameters.getTable(),
+ e);
+ }
+ }
+
@Override
public void close() {
if (dataClient != null) {
diff --git
a/seatunnel-connectors-v2/connector-google-bigtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/bigtable/source/BigtableSourceSplitEnumerator.java
b/seatunnel-connectors-v2/connector-google-bigtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/bigtable/source/BigtableSourceSplitEnumerator.java
index df72fda91d..b8d3f2b488 100644
---
a/seatunnel-connectors-v2/connector-google-bigtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/bigtable/source/BigtableSourceSplitEnumerator.java
+++
b/seatunnel-connectors-v2/connector-google-bigtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/bigtable/source/BigtableSourceSplitEnumerator.java
@@ -18,24 +18,37 @@
package org.apache.seatunnel.connectors.seatunnel.bigtable.source;
import org.apache.seatunnel.api.source.SourceSplitEnumerator;
+import
org.apache.seatunnel.connectors.seatunnel.bigtable.client.BigtableClient;
import
org.apache.seatunnel.connectors.seatunnel.bigtable.config.BigtableParameters;
+import com.google.cloud.bigtable.data.v2.models.KeyOffset;
import lombok.extern.slf4j.Slf4j;
import java.io.IOException;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashSet;
+import java.util.LinkedHashSet;
import java.util.List;
import java.util.Set;
import java.util.stream.Collectors;
/**
- * Enumerates {@link BigtableSourceSplit} instances for parallel reading.
+ * Enumerates {@link BigtableSourceSplit}s for parallel reading.
*
- * <p>Currently produces a single split covering the full table (or the
user-defined row-key range).
- * The split is assigned to whichever reader registers first. Future work can
partition by Bigtable
- * tablet boundaries using the Admin API.
+ * <p>Main workflow:
+ *
+ * <ol>
+ * <li>Use {@link BigtableClient#sampleRowKeys()} to split the key range
into approximately
+ * equal-size tablet boundaries.
+ * <li>Intersect each tablet range with the user-configured {@code
start_rowkey}/{@code
+ * end_rowkey}.
+ * <li>Assign splits to registered readers by hashing the split ID modulo
parallelism.
+ * </ol>
+ *
+ * <p>If split generation fails (sampling exception, empty samples, or empty
intersection), falls
+ * back to a single split covering the full user range so the job can still
proceed. Checkpoint
+ * persists both {@code assignedSplits} and {@code pendingSplits} (#11144).
*/
@Slf4j
public class BigtableSourceSplitEnumerator
@@ -47,24 +60,59 @@ public class BigtableSourceSplitEnumerator
private Set<BigtableSourceSplit> pendingSplits;
private boolean initialized = false;
+ /**
+ * Data API client used only for {@code sampleRowKeys}; injectable in unit
tests. Lazily created
+ * when split discovery first runs. The field itself is only published
under {@link #stateLock};
+ * construction may run outside the lock, but {@link #getBigtableClient()}
never publishes a
+ * client after {@link #closed} is set (see close/create race handling
there).
+ */
+ private BigtableClient bigtableClient;
+
+ /**
+ * Set under {@link #stateLock} by {@link #close()}. After this flag is
set, discovery must not
+ * publish a client, commit {@code pendingSplits}/{@code initialized}, or
fall back to a
+ * fabricated single split — shutdown wins over in-flight open().
+ */
+ private boolean closed = false;
+
/**
* Guards the shared assignment state ({@code assignedSplits}, {@code
pendingSplits}, {@code
- * initialized}) against concurrent enumerator callbacks. {@link
#initializePendingSplits()} and
- * {@link #assignSplit(int)} must only run while this lock is held.
+ * initialized}, {@code bigtableClient}, {@code closed}) against
concurrent enumerator
+ * callbacks. {@link #assignSplit(int)} must only run while this lock is
held. Blocking I/O
+ * ({@code sampleRowKeys()} and client construction) must run
<em>outside</em> this lock so it
+ * cannot stall {@link #snapshotState(long)} / checkpoint-barrier
processing.
*/
private final Object stateLock = new Object();
public BigtableSourceSplitEnumerator(
Context<BigtableSourceSplit> context, BigtableParameters
parameters) {
- this(context, parameters, null);
+ this(context, parameters, null, null);
}
public BigtableSourceSplitEnumerator(
Context<BigtableSourceSplit> context,
BigtableParameters parameters,
BigtableSourceState sourceState) {
+ this(context, parameters, sourceState, null);
+ }
+
+ /**
+ * Package-private constructor for injecting a {@link BigtableClient} in
unit tests to avoid
+ * real cloud connections.
+ *
+ * @param context engine split-assignment context
+ * @param parameters connection and scan parameters
+ * @param sourceState checkpoint recovery state; {@code null} on first
start
+ * @param bigtableClient pre-built client; lazily created on first split
if {@code null}
+ */
+ BigtableSourceSplitEnumerator(
+ Context<BigtableSourceSplit> context,
+ BigtableParameters parameters,
+ BigtableSourceState sourceState,
+ BigtableClient bigtableClient) {
this.context = context;
this.parameters = parameters;
+ this.bigtableClient = bigtableClient;
if (sourceState == null) {
this.assignedSplits = new HashSet<>();
this.pendingSplits = new HashSet<>();
@@ -83,7 +131,12 @@ public class BigtableSourceSplitEnumerator
@Override
public void open() {
- // State is fully initialized in the constructor; nothing to reset on
open().
+ // Discover splits before any reader registers. sampleRowKeys() is a
blocking RPC
+ // (plus lazy client construction); running it here avoids holding the
engine's
+ // enumeratorContext monitor used by triggerBarrier() during
receivedReader().
+ // Across Zeta/Flink/Spark, open() always completes before
registerReader() is called,
+ // so discovery is guaranteed to have finished by the time any reader
registers.
+ initializePendingSplits();
}
@Override
@@ -93,7 +146,13 @@ public class BigtableSourceSplitEnumerator
@Override
public void close() throws IOException {
- // Nothing to close – no persistent connection held here.
+ synchronized (stateLock) {
+ closed = true;
+ if (bigtableClient != null) {
+ bigtableClient.close();
+ bigtableClient = null;
+ }
+ }
}
@Override
@@ -110,13 +169,15 @@ public class BigtableSourceSplitEnumerator
@Override
public int currentUnassignedSplitSize() {
- return pendingSplits.size();
+ synchronized (stateLock) {
+ return pendingSplits.size();
+ }
}
@Override
public void registerReader(int subtaskId) {
+ // Discovery already ran in open(); assign under stateLock only.
synchronized (stateLock) {
- initializePendingSplits();
assignSplit(subtaskId);
}
}
@@ -134,36 +195,185 @@ public class BigtableSourceSplitEnumerator
@Override
public void handleSplitRequest(int subtaskId) {}
+ /**
+ * Discovers tablet splits once per job start. The blocking {@code
sampleRowKeys} RPC runs
+ * outside {@link #stateLock}; the lock is re-acquired only to mutate
{@code pendingSplits}.
+ *
+ * <p>If {@link #close()} wins a race against discovery, this method
returns without committing
+ * {@code pendingSplits} or {@code initialized}, so a fabricated fallback
split cannot be
+ * checkpointed after shutdown.
+ */
private void initializePendingSplits() {
- if (initialized) {
- return;
+ synchronized (stateLock) {
+ if (initialized || closed) {
+ return;
+ }
}
Set<BigtableSourceSplit> tableSplits = buildSplits();
- Set<String> existingIds =
- pendingSplits.stream()
- .map(BigtableSourceSplit::splitId)
- .collect(Collectors.toSet());
- existingIds.addAll(
- assignedSplits.stream()
- .map(BigtableSourceSplit::splitId)
- .collect(Collectors.toSet()));
- tableSplits.stream()
- .filter(s -> !existingIds.contains(s.splitId()))
- .forEach(pendingSplits::add);
- initialized = true;
+ synchronized (stateLock) {
+ // closed may have flipped while buildSplits() ran unlocked; do
not mutate after close.
+ if (initialized || closed) {
+ return;
+ }
+ Set<String> existingIds =
+ pendingSplits.stream()
+ .map(BigtableSourceSplit::splitId)
+ .collect(Collectors.toSet());
+ existingIds.addAll(
+ assignedSplits.stream()
+ .map(BigtableSourceSplit::splitId)
+ .collect(Collectors.toSet()));
+ tableSplits.stream()
+ .filter(s -> !existingIds.contains(s.splitId()))
+ .forEach(pendingSplits::add);
+ initialized = true;
+ }
}
/**
- * Builds the set of splits.
+ * Generates multiple splits via {@link BigtableClient#sampleRowKeys()};
falls back to a single
+ * split on failure.
*
- * <p>For now a single split spanning the requested row-key range is
produced. This is
- * sufficient for bounded batch reads. Parallel multi-split support can be
added later by
- * querying Bigtable tablet boundary information.
+ * <p>Sample keys are turned into half-open intervals {@code [prev,
current)} in lexicographic
+ * order, then intersected with the user range. If the last sample key is
non-empty (i.e. not
+ * the table-end sentinel), an extra segment {@code [lastSample, "")} is
appended to cover the
+ * table tail.
+ *
+ * <p>Performs blocking I/O (client construction and the sampleRowKeys
RPC). Callers must not
+ * hold {@link #stateLock} while invoking this method.
+ *
+ * @return discovered splits; a single split covering the user range when
sampling fails; an
+ * empty set when the enumerator was closed during discovery (caller
must not commit that)
*/
- private Set<BigtableSourceSplit> buildSplits() {
- String startKey = parameters.getStartRowkey() != null ?
parameters.getStartRowkey() : "";
- String endKey = parameters.getEndRowkey() != null ?
parameters.getEndRowkey() : "";
- return Collections.singleton(new BigtableSourceSplit(0, startKey,
endKey));
+ Set<BigtableSourceSplit> buildSplits() {
+ String userStart = parameters.getStartRowkey() != null ?
parameters.getStartRowkey() : "";
+ String userEnd = parameters.getEndRowkey() != null ?
parameters.getEndRowkey() : "";
+
+ List<KeyOffset> samples;
+ try {
+ samples = getBigtableClient().sampleRowKeys();
+ } catch (EnumeratorClosedException e) {
+ // close() raced discovery — not a Bigtable API failure; do not
log as sampleRowKeys
+ // failed and do not fabricate a whole-range fallback split.
+ log.info(
+ "Enumerator closed during split discovery for table [{}];
skipping discovery",
+ parameters.getTable());
+ return Collections.emptySet();
+ } catch (Exception e) {
+ if (isClosed()) {
+ log.info(
+ "Enumerator closed during split discovery for table
[{}]; skipping discovery",
+ parameters.getTable());
+ return Collections.emptySet();
+ }
+ log.warn(
+ "sampleRowKeys failed for table [{}], fallback to single
split",
+ parameters.getTable(),
+ e);
+ return Collections.singleton(new BigtableSourceSplit(0, userStart,
userEnd));
+ }
+
+ if (isClosed()) {
+ return Collections.emptySet();
+ }
+
+ if (samples == null || samples.isEmpty()) {
+ return Collections.singleton(new BigtableSourceSplit(0, userStart,
userEnd));
+ }
+
+ Set<BigtableSourceSplit> splits = new LinkedHashSet<>();
+ String rangeStart = "";
+ int index = 0;
+ for (KeyOffset sample : samples) {
+ String rangeEnd = keyToUtf8(sample);
+ index = addIntersectedSplit(splits, index, rangeStart, rangeEnd,
userStart, userEnd);
+ rangeStart = rangeEnd;
+ }
+
+ // If the last sample key is not empty the API did not emit a
table-end sentinel;
+ // append [lastSample, "") to avoid missing the table tail.
+ if (!rangeStart.isEmpty()) {
+ addIntersectedSplit(splits, index, rangeStart, "", userStart,
userEnd);
+ }
+
+ if (splits.isEmpty()) {
+ return Collections.singleton(new BigtableSourceSplit(0, userStart,
userEnd));
+ }
+ log.info(
+ "Enumerated {} Bigtable splits for table [{}]",
+ splits.size(),
+ parameters.getTable());
+ return splits;
+ }
+
+ /**
+ * Intersects a tablet range with the user range and adds the result to
the split set.
+ *
+ * @return {@code index + 1} if a split was added, otherwise the original
{@code index}
+ */
+ private static int addIntersectedSplit(
+ Set<BigtableSourceSplit> splits,
+ int index,
+ String rangeStart,
+ String rangeEnd,
+ String userStart,
+ String userEnd) {
+ String splitStart = maxStart(rangeStart, userStart);
+ String splitEnd = minEnd(rangeEnd, userEnd);
+ if (isValidRange(splitStart, splitEnd)) {
+ splits.add(new BigtableSourceSplit(index, splitStart, splitEnd));
+ return index + 1;
+ }
+ return index;
+ }
+
+ /**
+ * Converts a sample key to a UTF-8 string; null or empty ByteString
represents the table end.
+ */
+ private static String keyToUtf8(KeyOffset sample) {
+ if (sample == null || sample.getKey() == null) {
+ return "";
+ }
+ return sample.getKey().toStringUtf8();
+ }
+
+ /** Returns the lexicographically larger start key; empty string means
table-begin (minimum). */
+ private static String maxStart(String a, String b) {
+ if (a.isEmpty()) {
+ return b;
+ }
+ if (b.isEmpty()) {
+ return a;
+ }
+ return a.compareTo(b) >= 0 ? a : b;
+ }
+
+ /** Returns the lexicographically smaller end key; empty string means
table-end (maximum). */
+ private static String minEnd(String a, String b) {
+ if (a.isEmpty()) {
+ return b;
+ }
+ if (b.isEmpty()) {
+ return a;
+ }
+ return a.compareTo(b) <= 0 ? a : b;
+ }
+
+ /** Returns true if {@code [start, end)} is non-empty; empty end means
until table-end. */
+ private static boolean isValidRange(String start, String end) {
+ if (end.isEmpty()) {
+ return true;
+ }
+ if (start.isEmpty()) {
+ return true;
+ }
+ return start.compareTo(end) < 0;
+ }
+
+ private boolean isClosed() {
+ synchronized (stateLock) {
+ return closed;
+ }
}
private void assignSplit(int taskId) {
@@ -191,4 +401,53 @@ public class BigtableSourceSplitEnumerator
.collect(Collectors.joining(",")));
context.signalNoMoreSplits(taskId);
}
+
+ /**
+ * Returns the shared client, creating it outside {@link #stateLock} when
needed.
+ *
+ * <p>Construction is unlocked so a slow gRPC channel / credential load
does not stall
+ * checkpoint snapshotting. Before publishing the new instance, this
method re-checks under the
+ * lock: if another thread already published a client, or if {@link
#close()} has already set
+ * {@code closed}, the just-built instance is closed immediately and never
leaked.
+ *
+ * @throws EnumeratorClosedException if the enumerator is already closed
(or closes during
+ * construction), so callers can distinguish shutdown from a real
sampleRowKeys failure
+ */
+ private BigtableClient getBigtableClient() {
+ synchronized (stateLock) {
+ if (bigtableClient != null) {
+ return bigtableClient;
+ }
+ if (closed) {
+ throw new EnumeratorClosedException(
+ "BigtableSourceSplitEnumerator already closed; cannot
create client");
+ }
+ }
+ BigtableClient created = BigtableClient.createInstance(parameters);
+ synchronized (stateLock) {
+ if (closed) {
+ created.close();
+ throw new EnumeratorClosedException(
+ "BigtableSourceSplitEnumerator closed during client
creation");
+ }
+ if (bigtableClient == null) {
+ bigtableClient = created;
+ return created;
+ }
+ created.close();
+ return bigtableClient;
+ }
+ }
+
+ /**
+ * Thrown when discovery/client creation observes that {@link #close()}
has already run. Not a
+ * Bigtable API failure — callers must not treat it as {@code
sampleRowKeys} failure.
+ */
+ static final class EnumeratorClosedException extends IllegalStateException
{
+ private static final long serialVersionUID = 1L;
+
+ EnumeratorClosedException(String message) {
+ super(message);
+ }
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-google-bigtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/bigtable/source/BigtableSourceSplitEnumeratorTest.java
b/seatunnel-connectors-v2/connector-google-bigtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/bigtable/source/BigtableSourceSplitEnumeratorTest.java
index b0767c7252..c69e641672 100644
---
a/seatunnel-connectors-v2/connector-google-bigtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/bigtable/source/BigtableSourceSplitEnumeratorTest.java
+++
b/seatunnel-connectors-v2/connector-google-bigtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/bigtable/source/BigtableSourceSplitEnumeratorTest.java
@@ -21,20 +21,36 @@ import
org.apache.seatunnel.api.common.metrics.MetricsContext;
import org.apache.seatunnel.api.event.EventListener;
import org.apache.seatunnel.api.source.SourceEvent;
import org.apache.seatunnel.api.source.SourceSplitEnumerator;
+import
org.apache.seatunnel.connectors.seatunnel.bigtable.client.BigtableClient;
import
org.apache.seatunnel.connectors.seatunnel.bigtable.config.BigtableParameters;
+import
org.apache.seatunnel.connectors.seatunnel.bigtable.exception.BigtableConnectorErrorCode;
+import
org.apache.seatunnel.connectors.seatunnel.bigtable.exception.BigtableConnectorException;
import org.junit.jupiter.api.Test;
+import com.google.cloud.bigtable.data.v2.models.KeyOffset;
+import com.google.protobuf.ByteString;
+
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
public class BigtableSourceSplitEnumeratorTest {
@@ -61,15 +77,11 @@ public class BigtableSourceSplitEnumeratorTest {
@Test
void testReturnedSplitSurvivesCheckpointRestoreAndReassign() throws
Exception {
TestingContext context = new TestingContext(1);
- BigtableParameters parameters =
- BigtableParameters.builder()
- .projectId("test-project")
- .instanceId("test-instance")
- .table("test-table")
- .build();
+ BigtableParameters parameters = testParameters("", "");
+ BigtableClient client = mockClient(Collections.emptyList());
BigtableSourceSplitEnumerator enumerator =
- new BigtableSourceSplitEnumerator(context, parameters);
+ new BigtableSourceSplitEnumerator(context, parameters, null,
client);
enumerator.open();
context.registerReaderForTest(0);
@@ -86,7 +98,7 @@ public class BigtableSourceSplitEnumeratorTest {
assertTrue(checkpoint.getAssignedSplits().contains(split));
BigtableSourceSplitEnumerator restored =
- new BigtableSourceSplitEnumerator(context, parameters,
checkpoint);
+ new BigtableSourceSplitEnumerator(context, parameters,
checkpoint, client);
restored.open();
context.clearAssignments();
@@ -101,15 +113,11 @@ public class BigtableSourceSplitEnumeratorTest {
@Test
void testSnapshotStateDefensivelyCopiesPendingSplits() throws Exception {
TestingContext context = new TestingContext(1);
- BigtableParameters parameters =
- BigtableParameters.builder()
- .projectId("test-project")
- .instanceId("test-instance")
- .table("test-table")
- .build();
+ BigtableParameters parameters = testParameters("", "");
+ BigtableClient client = mockClient(Collections.emptyList());
BigtableSourceSplitEnumerator enumerator =
- new BigtableSourceSplitEnumerator(context, parameters);
+ new BigtableSourceSplitEnumerator(context, parameters, null,
client);
enumerator.open();
context.registerReaderForTest(0);
@@ -126,6 +134,275 @@ public class BigtableSourceSplitEnumeratorTest {
assertEquals(1, enumerator.currentUnassignedSplitSize());
}
+ @Test
+ void testSampleRowKeysProducesThreeSplits() {
+ BigtableSourceSplitEnumerator enumerator =
+ enumeratorWithSamples(testParameters("", ""),
Arrays.asList("m", "t", ""));
+
+ Set<BigtableSourceSplit> splits = enumerator.buildSplits();
+ assertEquals(3, splits.size());
+ assertEquals(
+ Arrays.asList("", "m", "t"),
+ splits.stream()
+ .map(BigtableSourceSplit::getStartRowKey)
+ .collect(Collectors.toList()));
+ assertEquals(
+ Arrays.asList("m", "t", ""),
+ splits.stream()
+ .map(BigtableSourceSplit::getEndRowKey)
+ .collect(Collectors.toList()));
+ }
+
+ @Test
+ void testUserRangeIntersectsSampledSplits() {
+ // Table tablets: ["","c") ["c","f") ["f",""); user range is [b,e)
+ BigtableSourceSplitEnumerator enumerator =
+ enumeratorWithSamples(testParameters("b", "e"),
Arrays.asList("c", "f", ""));
+
+ Set<BigtableSourceSplit> splits = enumerator.buildSplits();
+ assertEquals(2, splits.size());
+ List<String> starts =
+ splits.stream()
+ .map(BigtableSourceSplit::getStartRowKey)
+ .collect(Collectors.toList());
+ List<String> ends =
+
splits.stream().map(BigtableSourceSplit::getEndRowKey).collect(Collectors.toList());
+ assertEquals(Arrays.asList("b", "c"), starts);
+ assertEquals(Arrays.asList("c", "e"), ends);
+ }
+
+ @Test
+ void testSampleRowKeysFailureFallsBackToSingleSplit() {
+ BigtableClient client = mock(BigtableClient.class);
+ when(client.sampleRowKeys())
+ .thenThrow(
+ new BigtableConnectorException(
+ BigtableConnectorErrorCode.TABLE_QUERY_FAILED,
"boom"));
+ BigtableParameters parameters = testParameters("start", "end");
+ BigtableSourceSplitEnumerator enumerator =
+ new BigtableSourceSplitEnumerator(new TestingContext(1),
parameters, null, client);
+
+ Set<BigtableSourceSplit> splits = enumerator.buildSplits();
+ assertEquals(1, splits.size());
+ BigtableSourceSplit split = splits.iterator().next();
+ assertEquals("start", split.getStartRowKey());
+ assertEquals("end", split.getEndRowKey());
+ }
+
+ @Test
+ void testEmptySampleListFallsBackToSingleSplit() {
+ BigtableSourceSplitEnumerator enumerator =
+ enumeratorWithSamples(testParameters("s", "e"),
Collections.emptyList());
+
+ Set<BigtableSourceSplit> splits = enumerator.buildSplits();
+ assertEquals(1, splits.size());
+ BigtableSourceSplit split = splits.iterator().next();
+ assertEquals("s", split.getStartRowKey());
+ assertEquals("e", split.getEndRowKey());
+ }
+
+ @Test
+ void testEmptyIntersectionFallsBackToSingleSplit() {
+ // Inverted user range has no forward intersection with any sampled
interval
+ BigtableSourceSplitEnumerator enumerator =
+ enumeratorWithSamples(testParameters("z", "a"),
Arrays.asList("m", "t", ""));
+
+ Set<BigtableSourceSplit> splits = enumerator.buildSplits();
+ assertEquals(1, splits.size());
+ BigtableSourceSplit split = splits.iterator().next();
+ assertEquals("z", split.getStartRowKey());
+ assertEquals("a", split.getEndRowKey());
+ }
+
+ @Test
+ void testMissingTrailingEmptySampleStillCoversTableEnd() {
+ BigtableSourceSplitEnumerator enumerator =
+ enumeratorWithSamples(testParameters("", ""),
Arrays.asList("m", "t"));
+
+ Set<BigtableSourceSplit> splits = enumerator.buildSplits();
+ assertEquals(3, splits.size());
+ List<String> starts =
+ splits.stream()
+ .map(BigtableSourceSplit::getStartRowKey)
+ .collect(Collectors.toList());
+ List<String> ends =
+
splits.stream().map(BigtableSourceSplit::getEndRowKey).collect(Collectors.toList());
+ assertEquals(Arrays.asList("", "m", "t"), starts);
+ assertEquals(Arrays.asList("m", "t", ""), ends);
+ }
+
+ @Test
+ void testParallelismOneAssignsAllSampledSplits() throws Exception {
+ TestingContext context = new TestingContext(1);
+ BigtableParameters parameters = testParameters("", "");
+ BigtableClient client = mockClient(Arrays.asList("m", "t", ""));
+ BigtableSourceSplitEnumerator enumerator =
+ new BigtableSourceSplitEnumerator(context, parameters, null,
client);
+ enumerator.open();
+
+ context.registerReaderForTest(0);
+ enumerator.registerReader(0);
+
+ assertEquals(3, context.getAssignedSplitCount(0));
+ assertEquals(0, enumerator.currentUnassignedSplitSize());
+ }
+
+ /**
+ * Verifies that with parallelism N, each produced split is assigned to
exactly one reader and
+ * the assignment matches {@code hash(splitId) % N}.
+ *
+ * <p>Samples produce 3 splits (["","m"), ["m","t"), ["t","")). With
parallelism 3 every reader
+ * should receive at least one split, the union equals all splits, and no
split appears twice.
+ */
+ @Test
+ void testParallelismMultipleReadersEachGetDisjointHashedSplits() throws
Exception {
+ int parallelism = 3;
+ TestingContext context = new TestingContext(parallelism);
+ BigtableParameters parameters = testParameters("", "");
+ BigtableClient client = mockClient(Arrays.asList("m", "t", ""));
+ BigtableSourceSplitEnumerator enumerator =
+ new BigtableSourceSplitEnumerator(context, parameters, null,
client);
+ enumerator.open();
+
+ for (int i = 0; i < parallelism; i++) {
+ context.registerReaderForTest(i);
+ enumerator.registerReader(i);
+ }
+
+ // All 3 splits should have been assigned and nothing left pending.
+ assertEquals(0, enumerator.currentUnassignedSplitSize());
+
+ // Collect all assigned splits across all readers.
+ List<BigtableSourceSplit> allAssigned = new ArrayList<>();
+ for (int i = 0; i < parallelism; i++) {
+ allAssigned.addAll(context.getAssignedSplits(i));
+ }
+ assertEquals(3, allAssigned.size(), "Total splits across all readers
should be 3");
+
+ // No split should appear more than once (no duplicates).
+ Set<String> splitIds =
+
allAssigned.stream().map(BigtableSourceSplit::splitId).collect(Collectors.toSet());
+ assertEquals(3, splitIds.size(), "Each split must be assigned to
exactly one reader");
+
+ // Each split's owner must match hash(splitId) % parallelism.
+ for (int i = 0; i < parallelism; i++) {
+ for (BigtableSourceSplit split : context.getAssignedSplits(i)) {
+ int expected = (split.splitId().hashCode() &
Integer.MAX_VALUE) % parallelism;
+ assertEquals(
+ expected,
+ i,
+ "Split "
+ + split.splitId()
+ + " should belong to reader "
+ + expected
+ + " but was assigned to reader "
+ + i);
+ }
+ }
+ }
+
+ /**
+ * close() before open() must not fabricate a whole-range fallback split
or commit pending
+ * state. getBigtableClient() observes closed and discovery aborts without
a misleading
+ * sampleRowKeys-failure fallback.
+ */
+ @Test
+ void testCloseBeforeOpenDoesNotCommitDiscoveryState() throws Exception {
+ TestingContext context = new TestingContext(1);
+ BigtableClient client = mockClient(Arrays.asList("m", "t", ""));
+ BigtableSourceSplitEnumerator enumerator =
+ new BigtableSourceSplitEnumerator(context, testParameters("",
""), null, client);
+
+ enumerator.close();
+ enumerator.open();
+
+ assertEquals(0, enumerator.currentUnassignedSplitSize());
+ BigtableSourceState state = enumerator.snapshotState(1L);
+ assertTrue(state.getAssignedSplits().isEmpty());
+ assertTrue(state.getPendingSplits().isEmpty());
+ verify(client).close();
+ }
+
+ /**
+ * close() racing an in-flight sampleRowKeys() must not commit
pendingSplits/initialized after
+ * the RPC returns. Without the closed check in initializePendingSplits(),
empty samples would
+ * fall back to a fabricated whole-range split and persist it.
+ */
+ @Test
+ void testCloseDuringSampleRowKeysDoesNotCommitPendingSplits() throws
Exception {
+ CountDownLatch sampleStarted = new CountDownLatch(1);
+ CountDownLatch allowSampleToFinish = new CountDownLatch(1);
+
+ BigtableClient client = mock(BigtableClient.class);
+ when(client.sampleRowKeys())
+ .thenAnswer(
+ invocation -> {
+ sampleStarted.countDown();
+ assertTrue(allowSampleToFinish.await(10,
TimeUnit.SECONDS));
+ // Would normally trigger the single-split
empty-sample fallback.
+ return Collections.emptyList();
+ });
+
+ TestingContext context = new TestingContext(1);
+ BigtableSourceSplitEnumerator enumerator =
+ new BigtableSourceSplitEnumerator(context, testParameters("",
""), null, client);
+
+ ExecutorService pool = Executors.newSingleThreadExecutor();
+ Future<?> openFuture =
+ pool.submit(
+ () -> {
+ enumerator.open();
+ return null;
+ });
+
+ assertTrue(sampleStarted.await(10, TimeUnit.SECONDS));
+ enumerator.close();
+ allowSampleToFinish.countDown();
+ openFuture.get(10, TimeUnit.SECONDS);
+ pool.shutdown();
+ assertTrue(pool.awaitTermination(5, TimeUnit.SECONDS));
+
+ assertEquals(0, enumerator.currentUnassignedSplitSize());
+ BigtableSourceState state = enumerator.snapshotState(1L);
+ assertTrue(state.getAssignedSplits().isEmpty());
+ assertTrue(
+ state.getPendingSplits().isEmpty(),
+ "close() during discovery must not commit a fabricated
fallback split");
+ verify(client).close();
+ }
+
+ private static BigtableSourceSplitEnumerator enumeratorWithSamples(
+ BigtableParameters parameters, List<String> sampleKeys) {
+ return new BigtableSourceSplitEnumerator(
+ new TestingContext(1), parameters, null,
mockClient(sampleKeys));
+ }
+
+ private static BigtableParameters testParameters(String startRowkey,
String endRowkey) {
+ return BigtableParameters.builder()
+ .projectId("test-project")
+ .instanceId("test-instance")
+ .table("test-table")
+ .startRowkey(startRowkey.isEmpty() ? null : startRowkey)
+ .endRowkey(endRowkey.isEmpty() ? null : endRowkey)
+ .build();
+ }
+
+ private static BigtableClient mockClient(List<String> sampleKeys) {
+ BigtableClient client = mock(BigtableClient.class);
+ List<KeyOffset> samples = new ArrayList<>();
+ for (String key : sampleKeys) {
+ samples.add(keyOffset(key));
+ }
+ when(client.sampleRowKeys()).thenReturn(samples);
+ return client;
+ }
+
+ private static KeyOffset keyOffset(String key) {
+ KeyOffset offset = mock(KeyOffset.class);
+ when(offset.getKey()).thenReturn(key == null ? null :
ByteString.copyFromUtf8(key));
+ return offset;
+ }
+
private static class TestingContext
implements SourceSplitEnumerator.Context<BigtableSourceSplit> {
@@ -153,6 +430,11 @@ public class BigtableSourceSplitEnumeratorTest {
return assignments.getOrDefault(subtaskId,
Collections.emptyList()).size();
}
+ List<BigtableSourceSplit> getAssignedSplits(int subtaskId) {
+ return Collections.unmodifiableList(
+ assignments.getOrDefault(subtaskId,
Collections.emptyList()));
+ }
+
BigtableSourceSplit getLastAssignedSplit(int subtaskId) {
List<BigtableSourceSplit> splits = assignments.get(subtaskId);
return splits.get(splits.size() - 1);
diff --git
a/seatunnel-connectors-v2/connector-google-bigtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/bigtable/source/BigtableSourceStateRecoveryTest.java
b/seatunnel-connectors-v2/connector-google-bigtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/bigtable/source/BigtableSourceStateRecoveryTest.java
index e67b0b4ec5..0d9e294f8b 100644
---
a/seatunnel-connectors-v2/connector-google-bigtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/bigtable/source/BigtableSourceStateRecoveryTest.java
+++
b/seatunnel-connectors-v2/connector-google-bigtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/bigtable/source/BigtableSourceStateRecoveryTest.java
@@ -82,30 +82,53 @@ class BigtableSourceStateRecoveryTest {
.build();
/**
- * If the enumerator checkpoints before any reader registers, both sets
are empty. Restore must
- * still discover and assign the table split on the next {@link
- * BigtableSourceSplitEnumerator#registerReader(int)}.
+ * After eager open(), the first snapshot already contains discovered
pending splits (assigned
+ * still empty). Restoring that state must reassign the pending split
without re-running
+ * discovery.
*/
@Test
- void testEmptyEnumeratorCheckpointStillDiscoversSplitsOnRestore() throws
Exception {
+ void testCheckpointAfterEagerOpenRestoresPendingSplitsWithoutRediscovery()
throws Exception {
TestingContext context = new TestingContext(1);
- BigtableSourceSplitEnumerator enumerator =
- new BigtableSourceSplitEnumerator(context, PARAMETERS);
+ BigtableSourceSplitEnumerator enumerator = newEnumerator(context);
enumerator.open();
- BigtableSourceState emptyCheckpoint = enumerator.snapshotState(1L);
- assertTrue(emptyCheckpoint.getAssignedSplits().isEmpty());
- assertTrue(emptyCheckpoint.getPendingSplits().isEmpty());
+ BigtableSourceState checkpoint = enumerator.snapshotState(1L);
+ assertTrue(checkpoint.getAssignedSplits().isEmpty());
+ // open() eagerly discovers splits, so pendingSplits is non-empty at
first snapshot.
+ assertEquals(1, checkpoint.getPendingSplits().size());
+
+ BigtableSourceSplitEnumerator restored = restoreEnumerator(context,
checkpoint);
+ restored.open();
+
+ context.registerReaderForTest(0);
+ restored.registerReader(0);
+
+ assertEquals(1, context.getAssignedSplitCount(0));
+ assertEquals("bigtable_source_split_0",
context.getLastAssignedSplit(0).splitId());
+ }
+
+ /**
+ * Restoring a genuinely empty enumerator checkpoint (both assigned and
pending empty) must
+ * still discover splits on {@code open()}. This covers a barrier that
raced open()'s RPC and
+ * produced an empty snapshot, as well as restore from a pre-discovery
historical checkpoint.
+ */
+ @Test
+ void testTrulyEmptyCheckpointRestoreStillDiscoversSplitsOnOpen() throws
Exception {
+ TestingContext context = new TestingContext(1);
+ BigtableSourceState emptyState =
+ new BigtableSourceState(Collections.emptySet(),
Collections.emptySet());
- BigtableSourceSplitEnumerator restored =
- new BigtableSourceSplitEnumerator(context, PARAMETERS,
emptyCheckpoint);
+ BigtableSourceSplitEnumerator restored = restoreEnumerator(context,
emptyState);
+ // Constructor must treat empty-empty as not initialized so open()
rediscovers.
restored.open();
+ assertEquals(1, restored.currentUnassignedSplitSize());
context.registerReaderForTest(0);
restored.registerReader(0);
assertEquals(1, context.getAssignedSplitCount(0));
assertEquals("bigtable_source_split_0",
context.getLastAssignedSplit(0).splitId());
+ assertEquals(0, restored.currentUnassignedSplitSize());
}
/**
@@ -115,8 +138,7 @@ class BigtableSourceStateRecoveryTest {
@Test
void testReturnedSplitDroppedWhenCheckpointOmitsPending() throws Exception
{
TestingContext context = new TestingContext(1);
- BigtableSourceSplitEnumerator enumerator =
- new BigtableSourceSplitEnumerator(context, PARAMETERS);
+ BigtableSourceSplitEnumerator enumerator = newEnumerator(context);
enumerator.open();
context.registerReaderForTest(0);
@@ -133,8 +155,7 @@ class BigtableSourceStateRecoveryTest {
new HashSet<>(Collections.singleton(split)),
Collections.emptySet());
TestingContext restoreContext = new TestingContext(1);
- BigtableSourceSplitEnumerator restored =
- new BigtableSourceSplitEnumerator(restoreContext, PARAMETERS,
buggyCheckpoint);
+ BigtableSourceSplitEnumerator restored =
restoreEnumerator(restoreContext, buggyCheckpoint);
restored.open();
restoreContext.registerReaderForTest(0);
@@ -164,8 +185,7 @@ class BigtableSourceStateRecoveryTest {
assertTrue(legacyState.getPendingSplits().isEmpty());
TestingContext context = new TestingContext(1);
- BigtableSourceSplitEnumerator restored =
- new BigtableSourceSplitEnumerator(context, PARAMETERS,
legacyState);
+ BigtableSourceSplitEnumerator restored = restoreEnumerator(context,
legacyState);
restored.open();
context.registerReaderForTest(0);
@@ -191,8 +211,7 @@ class BigtableSourceStateRecoveryTest {
assertTrue(deserialized.getPendingSplits().isEmpty());
TestingContext context = new TestingContext(1);
- BigtableSourceSplitEnumerator restored =
- new BigtableSourceSplitEnumerator(context, PARAMETERS,
deserialized);
+ BigtableSourceSplitEnumerator restored = restoreEnumerator(context,
deserialized);
restored.open();
context.registerReaderForTest(0);
@@ -207,8 +226,7 @@ class BigtableSourceStateRecoveryTest {
@Test
void testReaderEnumeratorFailoverHandoff() throws Exception {
TestingContext enumContext = new TestingContext(1);
- BigtableSourceSplitEnumerator enumerator =
- new BigtableSourceSplitEnumerator(enumContext, PARAMETERS);
+ BigtableSourceSplitEnumerator enumerator = newEnumerator(enumContext);
enumerator.open();
enumContext.registerReaderForTest(0);
enumerator.registerReader(0);
@@ -236,7 +254,7 @@ class BigtableSourceStateRecoveryTest {
TestingContext restoreContext = new TestingContext(1);
BigtableSourceSplitEnumerator restoredEnum =
- new BigtableSourceSplitEnumerator(restoreContext, PARAMETERS,
enumCheckpoint);
+ restoreEnumerator(restoreContext, enumCheckpoint);
restoredEnum.open();
restoreContext.registerReaderForTest(0);
restoredEnum.registerReader(0);
@@ -320,8 +338,7 @@ class BigtableSourceStateRecoveryTest {
@Test
void testAddSplitsBackWithRegisteredReaderImmediatelyReassigns() throws
Exception {
TestingContext context = new TestingContext(1);
- BigtableSourceSplitEnumerator enumerator =
- new BigtableSourceSplitEnumerator(context, PARAMETERS);
+ BigtableSourceSplitEnumerator enumerator = newEnumerator(context);
enumerator.open();
context.registerReaderForTest(0);
@@ -339,6 +356,28 @@ class BigtableSourceStateRecoveryTest {
assertTrue(checkpoint.getAssignedSplits().contains(split));
}
+ private static BigtableSourceSplitEnumerator newEnumerator(
+ SourceSplitEnumerator.Context<BigtableSourceSplit> context) {
+ return new BigtableSourceSplitEnumerator(context, PARAMETERS, null,
emptySampleClient());
+ }
+
+ private static BigtableSourceSplitEnumerator restoreEnumerator(
+ SourceSplitEnumerator.Context<BigtableSourceSplit> context,
+ BigtableSourceState sourceState) {
+ return new BigtableSourceSplitEnumerator(
+ context, PARAMETERS, sourceState, emptySampleClient());
+ }
+
+ /**
+ * Empty samples force the single-split fallback so this class keeps its
one-split restore
+ * assertions.
+ */
+ private static BigtableClient emptySampleClient() {
+ BigtableClient client = mock(BigtableClient.class);
+ when(client.sampleRowKeys()).thenReturn(Collections.emptyList());
+ return client;
+ }
+
@SuppressWarnings("unchecked")
private static BigtableSourceReader createReaderWithMockedEmptyStream()
throws Exception {
BigtableClient mockClient = mock(BigtableClient.class);