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 6fd0857336 [Docs][Connector-V2] Improve Hbase Airtable Cassandra
Databend and Phoenix connector docs (#11799)
6fd0857336 is described below
commit 6fd0857336fd0206bc4fa0e96f0a4cdc94744377
Author: Daniel Carter <[email protected]>
AuthorDate: Fri Aug 14 17:06:02 2026 +0800
[Docs][Connector-V2] Improve Hbase Airtable Cassandra Databend and Phoenix
connector docs (#11799)
Co-authored-by: DanielCarter-stack
<[email protected]>
---
docs/en/connectors/sink/Airtable.md | 40 +++++++++++++++++++++++
docs/en/connectors/sink/Cassandra.md | 43 +++++++++++++++++++++++++
docs/en/connectors/sink/Databend.md | 36 +++++++++++++++++++++
docs/en/connectors/sink/Hbase.md | 20 ++++++++++++
docs/en/connectors/source/Airtable.md | 59 ++++++++++++++++++++++++++++++++++
docs/en/connectors/source/Cassandra.md | 21 +++++++++++-
docs/en/connectors/source/Databend.md | 25 +++++++++++++-
docs/en/connectors/source/Hbase.md | 47 +++++++++++++++++++++------
docs/en/connectors/source/Phoenix.md | 24 ++++++++++++++
docs/zh/connectors/sink/Airtable.md | 40 +++++++++++++++++++++++
docs/zh/connectors/sink/Cassandra.md | 40 +++++++++++++++++++++++
docs/zh/connectors/sink/Databend.md | 35 ++++++++++++++++++++
docs/zh/connectors/sink/Hbase.md | 20 ++++++++++++
docs/zh/connectors/source/Airtable.md | 58 +++++++++++++++++++++++++++++++++
docs/zh/connectors/source/Cassandra.md | 21 +++++++++++-
docs/zh/connectors/source/Databend.md | 22 ++++++++++++-
docs/zh/connectors/source/Hbase.md | 45 ++++++++++++++++++++------
docs/zh/connectors/source/Phoenix.md | 24 ++++++++++++++
18 files changed, 596 insertions(+), 24 deletions(-)
diff --git a/docs/en/connectors/sink/Airtable.md
b/docs/en/connectors/sink/Airtable.md
index 4cc3175475..c57d64a42b 100644
--- a/docs/en/connectors/sink/Airtable.md
+++ b/docs/en/connectors/sink/Airtable.md
@@ -161,6 +161,46 @@ sink {
}
```
+### Stream From Kafka To Airtable
+
+Combine a Kafka source with the Airtable sink to continuously push new events
into a tracking
+base. Keep `batch_size` at 10 to respect the Airtable request limit, and tune
+`request_interval_ms` if your topic produces bursts faster than 5 messages per
second.
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 60000
+}
+
+source {
+ Kafka {
+ bootstrap.servers = "kafka:9092"
+ topic = "orders.events"
+ format = "json"
+ schema = {
+ fields {
+ order_id = string
+ customer = string
+ amount = double
+ }
+ }
+ }
+}
+
+sink {
+ Airtable {
+ token = "patXXXXXXXX.XXXXXXXX"
+ base_id = "appXXXXXXXX"
+ table = "Orders"
+ typecast = true
+ batch_size = 10
+ request_interval_ms = 220
+ }
+}
+```
+
## Changelog
<ChangeLog />
\ No newline at end of file
diff --git a/docs/en/connectors/sink/Cassandra.md
b/docs/en/connectors/sink/Cassandra.md
index 662ec74d91..2be687384f 100644
--- a/docs/en/connectors/sink/Cassandra.md
+++ b/docs/en/connectors/sink/Cassandra.md
@@ -170,6 +170,49 @@ sink {
}
```
+### Stream MySQL CDC Events Into Cassandra
+
+Pipe MySQL CDC events through a Cassandra sink by mapping the CDC row kinds to
the table
+columns. Cassandra uses `INSERT` semantics per row, so a CDC `DELETE` is
expressed by writing
+a tombstone column (`is_deleted`) and filtering it out downstream if needed:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 10000
+}
+
+source {
+ MySQL-CDC {
+ base-url = "jdbc:mysql://mysql:3306/test"
+ username = "root"
+ password = "mysqlpw"
+ table-names = ["test.orders"]
+ }
+}
+
+sink {
+ Cassandra {
+ host = "cassandra1:9042,cassandra2:9042"
+ keyspace = "test"
+ table = "orders"
+ fields = ["id", "order_id", "customer", "amount", "is_deleted"]
+ consistency_level = "LOCAL_QUORUM"
+ batch_size = 2000
+ batch_type = "UNLOGGED"
+ async_write = true
+ }
+}
+```
+
+> **Note:** The `is_deleted` column shown above is not produced by MySQL-CDC
and is not
+> derived from `RowKind` by the Cassandra sink. You must populate it yourself
— either by
+> carrying an `is_deleted` column in the upstream MySQL table, or by adding a
Transform-V2
+> (for example `sql` / `replace`) between the source and the sink that
synthesizes it from
+> the CDC `RowKind`. Without that, `DELETE` events will be written back to
Cassandra as a
+> regular upsert.
+
## Changelog
<ChangeLog />
diff --git a/docs/en/connectors/sink/Databend.md
b/docs/en/connectors/sink/Databend.md
index 1fe778b55b..ee9edaa1f0 100644
--- a/docs/en/connectors/sink/Databend.md
+++ b/docs/en/connectors/sink/Databend.md
@@ -235,6 +235,42 @@ sink {
}
```
+### Stream MySQL CDC To Databend In Streaming Mode
+
+The same CDC settings also work in streaming jobs. The following example pipes
MySQL CDC
+events into Databend continuously. Keep `batch_size` small in streaming CDC
jobs so that each
+checkpoint reflects the latest writes:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 10000
+}
+
+source {
+ MySQL-CDC {
+ base-url = "jdbc:mysql://mysql:3306/test"
+ username = "root"
+ password = "mysqlpw"
+ table-names = ["test.orders"]
+ }
+}
+
+sink {
+ Databend {
+ url = "jdbc:databend://databend:8000/default?ssl=false"
+ username = "root"
+ password = ""
+ database = "default"
+ table = "orders"
+ batch_size = 500
+ conflict_key = "id"
+ enable_delete = true
+ }
+}
+```
+
## Related Links
- [Databend Official Website](https://databend.rs/)
diff --git a/docs/en/connectors/sink/Hbase.md b/docs/en/connectors/sink/Hbase.md
index 868070a7d1..3c0ed011ba 100644
--- a/docs/en/connectors/sink/Hbase.md
+++ b/docs/en/connectors/sink/Hbase.md
@@ -184,6 +184,26 @@ sink {
}
```
+### Recreate the Target Table Before Writing
+
+Use `schema_save_mode = "RECREATE_SCHEMA"` when the job is allowed to drop and
re-create the
+target table on every run. Combine it with `data_save_mode = "DROP_DATA"` to
wipe any previous
+rows, or keep `APPEND_DATA` if you only want the table structure refreshed.
+
+```hocon
+sink {
+ Hbase {
+ zookeeper_quorum = "hbase_e2e:2181"
+ table = "seatunnel_test_with_recreate_schema"
+ rowkey_column = ["name"]
+ family_name {
+ all_columns = info
+ }
+ schema_save_mode = "RECREATE_SCHEMA"
+ }
+}
+```
+
## Kerberos Example
Note:
diff --git a/docs/en/connectors/source/Airtable.md
b/docs/en/connectors/source/Airtable.md
index 63c492fd0e..38ecef27df 100644
--- a/docs/en/connectors/source/Airtable.md
+++ b/docs/en/connectors/source/Airtable.md
@@ -174,6 +174,65 @@ source {
}
```
+### Restrict To A View
+
+Use `view` to read only records that are visible in a specific view. Combine
it with `fields` to
+project only the columns the view exposes:
+
+```hocon
+source {
+ Airtable {
+ token = "patXXXXXXXX.XXXXXXXX"
+ base_id = "appXXXXXXXX"
+ table = "Shipments"
+ view = "Pending shipments"
+ fields = ["Name", "Status", "Weight"]
+ format = "json"
+ content_field = "$.records[*].fields"
+ schema = {
+ fields {
+ Name = string
+ Status = string
+ Weight = float
+ }
+ }
+ }
+}
+```
+
+### Run An Incremental Batch Read
+
+The Airtable source only supports `BATCH` jobs (it rejects non-batch modes).
To consume
+newly added rows across runs, pin `filter_by_formula` together with a `sort`
that orders rows by
+`createdTime`, and re-inject the last-seen `createdTime` watermark between
runs:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ Airtable {
+ token = "patXXXXXXXX.XXXXXXXX"
+ base_id = "appXXXXXXXX"
+ table = "Shipments"
+ format = "json"
+ content_field = "$.records[*].fields"
+ filter_by_formula = "IS_AFTER({CreatedAt}, '2026-01-01T00:00:00.000Z')"
+ sort = "[{\"field\":\"CreatedAt\",\"direction\":\"asc\"}]"
+ page_size = 100
+ schema = {
+ fields {
+ Name = string
+ Status = string
+ CreatedAt = string
+ }
+ }
+ }
+}
+```
+
## Changelog
<ChangeLog />
\ No newline at end of file
diff --git a/docs/en/connectors/source/Cassandra.md
b/docs/en/connectors/source/Cassandra.md
index 0b9af9238d..5b240cb1fb 100644
--- a/docs/en/connectors/source/Cassandra.md
+++ b/docs/en/connectors/source/Cassandra.md
@@ -105,7 +105,7 @@ duplicate table names during startup.
Example entry:
-```
+```hocon
{
cql = "SELECT id, name FROM keyspace.table1"
}
@@ -210,6 +210,25 @@ sink {
}
```
+### Read With A Stricter Consistency Level
+
+Use `consistency_level = "QUORUM"` when the read result must satisfy the
configured replication
+factor. Combine it with `datacenter` so the driver talks to the right local
coordinator:
+
+```hocon
+source {
+ Cassandra {
+ host = "cassandra1:9042,cassandra2:9042"
+ username = "cassandra"
+ password = "cassandra"
+ datacenter = "datacenter1"
+ keyspace = "test"
+ consistency_level = "QUORUM"
+ cql = "SELECT id, name, score FROM test.accounts"
+ }
+}
+```
+
## Changelog
<ChangeLog />
diff --git a/docs/en/connectors/source/Databend.md
b/docs/en/connectors/source/Databend.md
index ef50a692ce..501ac0b01c 100644
--- a/docs/en/connectors/source/Databend.md
+++ b/docs/en/connectors/source/Databend.md
@@ -21,7 +21,14 @@ import ChangeLog from '../changelog/connector-databend.md';
## Description
-A source connector for reading data from Databend.
+A source connector for reading data from [Databend](https://databend.rs/)
using the Databend JDBC
+driver. You can read a single table with `database` + `table`, run an ad-hoc
query with `query`,
+or supply an explicit statement through `sql`. The connector executes the
query in batch mode and
+returns each result row as a SeaTunnel row.
+
+The connector supports column projection through standard SQL and exposes a
few JDBC tuning
+options such as `fetch_size` and `ssl`. It does not currently read multiple
tables in a single
+source block; configure one Databend source per table.
## Dependencies
@@ -135,6 +142,22 @@ source {
}
```
+### Filter On A Computed Column
+
+Use any expression supported by Databend in `query` to project and filter rows
before they reach
+SeaTunnel:
+
+```hocon
+source {
+ Databend {
+ url = "jdbc:databend://localhost:8000"
+ username = "root"
+ password = ""
+ query = "SELECT id, name, age FROM default.users WHERE age >= 18 AND
starts_with(name, 'A') ORDER BY id"
+ }
+}
+```
+
## Related Links
- [Databend Official Website](https://databend.rs/)
diff --git a/docs/en/connectors/source/Hbase.md
b/docs/en/connectors/source/Hbase.md
index bf4f13b67f..612bbde0bb 100644
--- a/docs/en/connectors/source/Hbase.md
+++ b/docs/en/connectors/source/Hbase.md
@@ -15,6 +15,9 @@ import ChangeLog from '../changelog/connector-hbase.md';
Reads data from Apache HBase tables. The source supports normal scans, row key
range scans, timestamp
range scans, binary row keys, custom namespaces, and parallel batch reading.
+The source runs in batch mode and reads the current state of the scanned
range. It is not a CDC
+source: changes that happen after the scan starts are not delivered.
+
## Key Features
- [x] [batch](../../introduction/concepts/connector-v2-features.md)
@@ -125,14 +128,14 @@ Common parameters for Source plugins, refer to [Common
Source Options](../common
### Read by Row Key and Time Range
-```bash
+```hocon
source {
Hbase {
- zookeeper_quorum = "hadoop001:2181,hadoop002:2181,hadoop003:2181"
- table = "seatunnel_test"
- caching = 1000
- batch = 100
- cache_blocks = false
+ zookeeper_quorum = "hadoop001:2181,hadoop002:2181,hadoop003:2181"
+ table = "seatunnel_test"
+ caching = 1000
+ batch = 100
+ cache_blocks = false
is_binary_rowkey = false
start_rowkey = "B"
end_rowkey = "C"
@@ -140,16 +143,16 @@ source {
end_timestamp = 1700003600000
schema = {
columns = [
- {
- name = "rowkey"
- type = string
+ {
+ name = "rowkey"
+ type = string
},
{
name = "columnFamily1:column1"
type = boolean
},
{
- name = "columnFamily1:column2"
+ name = "columnFamily1:column2"
type = double
},
{
@@ -179,6 +182,30 @@ source {
}
```
+### Read with a Binary Row Key
+
+```hocon
+source {
+ Hbase {
+ zookeeper_quorum = "hbase_e2e:2181"
+ table = "binary_rowkey_table"
+ is_binary_rowkey = true
+ caching = 500
+ batch = 100
+ schema = {
+ columns = [
+ { name = rowkey, type = bytes },
+ { name = "info:name", type = string },
+ { name = "info:score", type = double }
+ ]
+ }
+ }
+}
+```
+
+When `is_binary_rowkey = true`, declare the row key column as `bytes` in the
schema and let the
+downstream transform handle decoding.
+
## Kerberos Example
Note:
diff --git a/docs/en/connectors/source/Phoenix.md
b/docs/en/connectors/source/Phoenix.md
index 89b09c0994..f68bb7b609 100644
--- a/docs/en/connectors/source/Phoenix.md
+++ b/docs/en/connectors/source/Phoenix.md
@@ -110,6 +110,30 @@ sink {
}
```
+### Project Specific Columns With A Predicate
+
+Combine column projection with a `WHERE` clause to narrow the rows that flow
downstream. This
+example only fetches rows whose `name` starts with `A`:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ Jdbc {
+ driver = org.apache.phoenix.jdbc.PhoenixDriver
+ url = "jdbc:phoenix:localhost:2182/hbase"
+ query = "select name, score from test.source where name like 'A%'"
+ }
+}
+
+sink {
+ Console {}
+}
+```
+
## Changelog
<ChangeLog />
\ No newline at end of file
diff --git a/docs/zh/connectors/sink/Airtable.md
b/docs/zh/connectors/sink/Airtable.md
index 6778c1bcae..25984435cb 100644
--- a/docs/zh/connectors/sink/Airtable.md
+++ b/docs/zh/connectors/sink/Airtable.md
@@ -161,6 +161,46 @@ sink {
}
```
+### 从 Kafka 流式写入 Airtable
+
+将 Kafka 源与 Airtable sink 组合,可以把订单等事件持续推送到运营跟踪表。
+`batch_size` 保持在 10 以满足 Airtable 的请求上限;当 Topic 的突发流量超过每秒 5 条时,
+适当上调 `request_interval_ms`。
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 60000
+}
+
+source {
+ Kafka {
+ bootstrap.servers = "kafka:9092"
+ topic = "orders.events"
+ format = "json"
+ schema = {
+ fields {
+ order_id = string
+ customer = string
+ amount = double
+ }
+ }
+ }
+}
+
+sink {
+ Airtable {
+ token = "patXXXXXXXX.XXXXXXXX"
+ base_id = "appXXXXXXXX"
+ table = "Orders"
+ typecast = true
+ batch_size = 10
+ request_interval_ms = 220
+ }
+}
+```
+
## 变更日志
<ChangeLog />
\ No newline at end of file
diff --git a/docs/zh/connectors/sink/Cassandra.md
b/docs/zh/connectors/sink/Cassandra.md
index 48affe8fb9..02c943892f 100644
--- a/docs/zh/connectors/sink/Cassandra.md
+++ b/docs/zh/connectors/sink/Cassandra.md
@@ -163,6 +163,46 @@ sink {
}
```
+### 将 MySQL CDC 事件流式写入 Cassandra
+
+把 MySQL CDC 事件通过 Cassandra sink 写入下游,需要把 CDC 的行类型映射成表字段。Cassandra
+按行执行 `INSERT`,CDC 的 `DELETE` 可以通过写入 `is_deleted` 字段并在下游过滤掉来实现:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 10000
+}
+
+source {
+ MySQL-CDC {
+ base-url = "jdbc:mysql://mysql:3306/test"
+ username = "root"
+ password = "mysqlpw"
+ table-names = ["test.orders"]
+ }
+}
+
+sink {
+ Cassandra {
+ host = "cassandra1:9042,cassandra2:9042"
+ keyspace = "test"
+ table = "orders"
+ fields = ["id", "order_id", "customer", "amount", "is_deleted"]
+ consistency_level = "LOCAL_QUORUM"
+ batch_size = 2000
+ batch_type = "UNLOGGED"
+ async_write = true
+ }
+}
+```
+
+> **注意**:上面示例中的 `is_deleted` 字段不会由 MySQL-CDC 自动产出,也不会被 Cassandra
+> sink 根据 `RowKind` 推导。你需要自行提供——既可以让上游 MySQL 表本身带有 `is_deleted`
+> 列,也可以在 source 和 sink 之间增加一个 Transform-V2(例如 `sql`、`replace`)从 CDC 的
+> `RowKind` 合成该字段。否则 `DELETE` 事件会被当作普通 upsert 写回 Cassandra。
+
## 变更日志
<ChangeLog />
diff --git a/docs/zh/connectors/sink/Databend.md
b/docs/zh/connectors/sink/Databend.md
index 7f3b3423c4..188f087a36 100644
--- a/docs/zh/connectors/sink/Databend.md
+++ b/docs/zh/connectors/sink/Databend.md
@@ -233,6 +233,41 @@ sink {
}
```
+### 将 MySQL CDC 流式写入 Databend
+
+同一套 CDC 参数同样适用于流式任务。下面的示例将 MySQL CDC 事件持续写入 Databend。流式
+CDC 场景下建议把 `batch_size` 调小一些,让每个 checkpoint 都能反映最新写入。
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 10000
+}
+
+source {
+ MySQL-CDC {
+ base-url = "jdbc:mysql://mysql:3306/test"
+ username = "root"
+ password = "mysqlpw"
+ table-names = ["test.orders"]
+ }
+}
+
+sink {
+ Databend {
+ url = "jdbc:databend://databend:8000/default?ssl=false"
+ username = "root"
+ password = ""
+ database = "default"
+ table = "orders"
+ batch_size = 500
+ conflict_key = "id"
+ enable_delete = true
+ }
+}
+```
+
## 相关链接
- [Databend 官方网站](https://databend.rs/)
diff --git a/docs/zh/connectors/sink/Hbase.md b/docs/zh/connectors/sink/Hbase.md
index 2ee8220aef..4644fb6926 100644
--- a/docs/zh/connectors/sink/Hbase.md
+++ b/docs/zh/connectors/sink/Hbase.md
@@ -178,6 +178,26 @@ sink {
}
```
+### 每次写入前重建目标表
+
+当任务允许在每次运行时删除并重建目标表时,可以使用 `schema_save_mode = "RECREATE_SCHEMA"`。
+配合 `data_save_mode = "DROP_DATA"` 可以同时清空已有数据;如果只想刷新表结构,保留
+`APPEND_DATA` 即可。
+
+```hocon
+sink {
+ Hbase {
+ zookeeper_quorum = "hbase_e2e:2181"
+ table = "seatunnel_test_with_recreate_schema"
+ rowkey_column = ["name"]
+ family_name {
+ all_columns = info
+ }
+ schema_save_mode = "RECREATE_SCHEMA"
+ }
+}
+```
+
## Kerberos 示例
备注:
diff --git a/docs/zh/connectors/source/Airtable.md
b/docs/zh/connectors/source/Airtable.md
index be344a3333..b6e29088e4 100644
--- a/docs/zh/connectors/source/Airtable.md
+++ b/docs/zh/connectors/source/Airtable.md
@@ -174,6 +174,64 @@ source {
}
```
+### 仅读取某个视图中的记录
+
+使用 `view` 只读取指定视图中可见的记录,配合 `fields` 限制返回的列:
+
+```hocon
+source {
+ Airtable {
+ token = "patXXXXXXXX.XXXXXXXX"
+ base_id = "appXXXXXXXX"
+ table = "Shipments"
+ view = "Pending shipments"
+ fields = ["Name", "Status", "Weight"]
+ format = "json"
+ content_field = "$.records[*].fields"
+ schema = {
+ fields {
+ Name = string
+ Status = string
+ Weight = float
+ }
+ }
+ }
+}
+```
+
+### 按增量批次消费新增记录
+
+Airtable source 只支持 `BATCH` 作业(其它模式在初始化时会被拒绝)。要在多次运行之间持续
+消费新增行,可以固定 `filter_by_formula` 并配合 `sort` 按 `createdTime` 升序排列,同时把
+上一次运行读到的最大 `CreatedAt` 写回公式:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ Airtable {
+ token = "patXXXXXXXX.XXXXXXXX"
+ base_id = "appXXXXXXXX"
+ table = "Shipments"
+ format = "json"
+ content_field = "$.records[*].fields"
+ filter_by_formula = "IS_AFTER({CreatedAt}, '2026-01-01T00:00:00.000Z')"
+ sort = "[{\"field\":\"CreatedAt\",\"direction\":\"asc\"}]"
+ page_size = 100
+ schema = {
+ fields {
+ Name = string
+ Status = string
+ CreatedAt = string
+ }
+ }
+ }
+}
+```
+
## 变更日志
<ChangeLog />
\ No newline at end of file
diff --git a/docs/zh/connectors/source/Cassandra.md
b/docs/zh/connectors/source/Cassandra.md
index 687a0bf1a0..91ce6ac5a1 100644
--- a/docs/zh/connectors/source/Cassandra.md
+++ b/docs/zh/connectors/source/Cassandra.md
@@ -100,7 +100,7 @@ Cassandra source 支持两种读取方式:
示例条目:
-```
+```hocon
{
cql = "SELECT id, name FROM keyspace.table1"
}
@@ -203,6 +203,25 @@ sink {
}
```
+### 提高读取一致性级别
+
+当读取结果必须满足配置的副本因子时,使用 `consistency_level = "QUORUM"`,并配合
+`datacenter` 让 Driver 连接到正确的本地协调节点:
+
+```hocon
+source {
+ Cassandra {
+ host = "cassandra1:9042,cassandra2:9042"
+ username = "cassandra"
+ password = "cassandra"
+ datacenter = "datacenter1"
+ keyspace = "test"
+ consistency_level = "QUORUM"
+ cql = "SELECT id, name, score FROM test.accounts"
+ }
+}
+```
+
## 变更日志
<ChangeLog />
diff --git a/docs/zh/connectors/source/Databend.md
b/docs/zh/connectors/source/Databend.md
index e1642b9be8..39ec0fcf52 100644
--- a/docs/zh/connectors/source/Databend.md
+++ b/docs/zh/connectors/source/Databend.md
@@ -20,7 +20,12 @@ import ChangeLog from '../changelog/connector-databend.md';
## 描述
-用于从 Databend 读取数据的源连接器。
+通过 Databend JDBC 驱动从 [Databend](https://databend.rs/) 读取数据的源连接器。可以使用
+`database` + `table` 读取单张表,使用 `query` 执行一次性查询,也可以通过 `sql` 提供完整的
+SQL 语句。连接器在批处理模式下执行查询,并把每一行结果转换为 SeaTunnel 行。
+
+连接器支持标准 SQL 的列投影,并提供 `fetch_size`、`ssl` 等 JDBC 调优选项。当前不支持在同
+一个 source 块中读取多张表,需要为每张表配置一个 Databend source。
## 依赖
@@ -133,6 +138,21 @@ source {
}
```
+### 在查询中过滤和投影
+
+可以直接在 `query` 里使用 Databend 支持的任意表达式,提前把不需要的列过滤掉:
+
+```hocon
+source {
+ Databend {
+ url = "jdbc:databend://localhost:8000"
+ username = "root"
+ password = ""
+ query = "SELECT id, name, age FROM default.users WHERE age >= 18 AND
starts_with(name, 'A') ORDER BY id"
+ }
+}
+```
+
## 相关链接
- [Databend 官方网站](https://databend.rs/)
diff --git a/docs/zh/connectors/source/Hbase.md
b/docs/zh/connectors/source/Hbase.md
index 85a6f04a3e..8b2d0c73b6 100644
--- a/docs/zh/connectors/source/Hbase.md
+++ b/docs/zh/connectors/source/Hbase.md
@@ -14,6 +14,8 @@ import ChangeLog from '../changelog/connector-hbase.md';
从 Apache HBase 表读取数据。支持普通扫描、行键范围扫描、时间戳范围扫描、二进制行键、自定义命名空间和并行批读取。
+本连接器以批处理模式运行,读取的是扫描范围在任务启动时刻的快照状态,并非 CDC 源;扫描启动后发生的新写入不会进入结果。
+
## 主要特性
- [x] [批处理](../../introduction/concepts/connector-v2-features.md)
@@ -124,14 +126,14 @@ Source 插件常用参数,具体请参考 [Source 常用选项](../common-opti
### 按行键和时间范围读取
-```bash
+```hocon
source {
Hbase {
- zookeeper_quorum = "hadoop001:2181,hadoop002:2181,hadoop003:2181"
- table = "seatunnel_test"
- caching = 1000
- batch = 100
- cache_blocks = false
+ zookeeper_quorum = "hadoop001:2181,hadoop002:2181,hadoop003:2181"
+ table = "seatunnel_test"
+ caching = 1000
+ batch = 100
+ cache_blocks = false
is_binary_rowkey = false
start_rowkey = "B"
end_rowkey = "C"
@@ -139,16 +141,16 @@ source {
end_timestamp = 1700003600000
schema = {
columns = [
- {
- name = "rowkey"
- type = string
+ {
+ name = "rowkey"
+ type = string
},
{
name = "columnFamily1:column1"
type = boolean
},
{
- name = "columnFamily1:column2"
+ name = "columnFamily1:column2"
type = double
},
{
@@ -178,6 +180,29 @@ source {
}
```
+### 读取二进制行键
+
+```hocon
+source {
+ Hbase {
+ zookeeper_quorum = "hbase_e2e:2181"
+ table = "binary_rowkey_table"
+ is_binary_rowkey = true
+ caching = 500
+ batch = 100
+ schema = {
+ columns = [
+ { name = rowkey, type = bytes },
+ { name = "info:name", type = string },
+ { name = "info:score", type = double }
+ ]
+ }
+ }
+}
+```
+
+当 `is_binary_rowkey = true` 时,请在 `schema` 中把行键列声明为 `bytes`,并在后续的 transform
节点里自行解码。
+
## Kerberos 示例
备注:
diff --git a/docs/zh/connectors/source/Phoenix.md
b/docs/zh/connectors/source/Phoenix.md
index 2d383e3a39..d1457d5c8a 100644
--- a/docs/zh/connectors/source/Phoenix.md
+++ b/docs/zh/connectors/source/Phoenix.md
@@ -110,6 +110,30 @@ sink {
}
```
+### 列投影并按条件过滤
+
+可以结合列投影和 `WHERE` 条件提前把不需要的行过滤掉。下面的示例只读取 `name` 以 `A`
+开头的行:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ Jdbc {
+ driver = org.apache.phoenix.jdbc.PhoenixDriver
+ url = "jdbc:phoenix:localhost:2182/hbase"
+ query = "select name, score from test.source where name like 'A%'"
+ }
+}
+
+sink {
+ Console {}
+}
+```
+
## 变更日志
<ChangeLog />
\ No newline at end of file