Skip to main content
Version: 3.0.0

Kudu

Kudu source connector

Support Kudu Version​

  • 1.11.1/1.12.0/1.13.0/1.14.0/1.15.0

Support Those Engines​

Spark
Flink
SeaTunnel Zeta

Key features​

Description​

Used to read data from Kudu.

The tested kudu version is 1.11.1.

Data Type Mapping​

kudu Data TypeSeaTunnel Data Type
BOOLBOOLEAN
INT8
INT16
INT32
INT
INT64BIGINT
DECIMALDECIMAL
FLOATFLOAT
DOUBLEDOUBLE
STRINGSTRING
UNIXTIME_MICROSTIMESTAMP
BINARYBYTES

Source Options​

NameTypeRequiredDefaultDescription
kudu_mastersStringYes-Kudu master address. Separated by ',',such as '192.168.88.110:7051'.
table_nameStringYes, if table_list is not configured-The name of the Kudu table. This option is mutually exclusive with table_list.
client_worker_countIntNo2 * Runtime.getRuntime().availableProcessors()Kudu worker count. Default value is twice the current number of cpu cores.
client_default_operation_timeout_msLongNo30000Kudu normal operation time out.
client_default_admin_operation_timeout_msLongNo30000Kudu admin operation time out.
enable_kerberosBoolNofalseKerberos principal enable.
kerberos_principalStringYes, when enable_kerberos = true-Kerberos principal used by the Kudu client. The keytab must be available on every worker node.
kerberos_keytabStringYes, when enable_kerberos = true-Kerberos keytab path used by the Kudu client. The file must be available on every worker node.
kerberos_krb5confStringNo-Kerberos krb5 conf. Note that all zeta nodes require have this file.
scan_token_query_timeoutLongNo30000The timeout for connecting scan token. If not set, it will be the same as operationTimeout.
scan_token_batch_size_bytesIntNo1024 * 1024Kudu scan bytes. The maximum number of bytes read at a time, the default is 1MB.
use_regexBoolNofalseControl regular expression matching for table_name. When set to true, the table_name will be treated as a regular expression pattern and can match multiple tables. When set to false or not specified, the table_name will be treated as an exact table name (no regex matching).
filterStringNo-Kudu scan filter expressions,example id > 100 AND id < 200.
schemaMapNo-SeaTunnel Schema. For more details, please refer to Schema Feature.
table_listArrayNo-The list of tables to read. Use this option instead of table_name, for example: table_list = [{ table_name = "kudu_source_table_1"},{ table_name = "kudu_source_table_2"}] . You can also configure use_regex = true inside each entry to enable regex matching for that table_name.
common-optionsNo-Source plugin common parameters, please refer to Source Common Options for details.

Option Notes​

  • Configure exactly one of table_name and table_list.
  • filter is pushed down to Kudu scans and can use Kudu predicate expressions such as id >= 1 AND id <= 2.
  • use_regex = true treats table_name as a Java regular expression. This can be used either at the top level or inside each table_list item.
  • When enable_kerberos = true, both kerberos_principal and kerberos_keytab are required.

Task Example​

Simple​

The following example is for a Kudu table named "kudu_source_table", The goal is to print the data from this table on the console and write kudu table "kudu_sink_table"

# Defining the runtime environment
env {
parallelism = 2
job.mode = "BATCH"
}

source {
# This is a example source plugin **only for test and demonstrate the feature source plugin**
kudu {
kudu_masters = "kudu-master:7051"
table_name = "kudu_source_table"
plugin_output = "kudu"
enable_kerberos = true
kerberos_principal = "xx@xx.COM"
kerberos_keytab = "xx.keytab"
}
}

transform {
}

sink {
console {
plugin_input = "kudu"
}

kudu {
plugin_input = "kudu"
kudu_masters = "kudu-master:7051"
table_name = "kudu_sink_table"
enable_kerberos = true
kerberos_principal = "xx@xx.COM"
kerberos_keytab = "xx.keytab"
}
}

Multiple Table​

env {
# You can set engine configuration here
parallelism = 1
job.mode = "STREAMING"
checkpoint.interval = 5000
}

source {
# This is a example source plugin **only for test and demonstrate the feature source plugin**
kudu{
kudu_masters = "kudu-master:7051"
table_list = [
{
table_name = "kudu_source_table_1"
},{
table_name = "kudu_source_table_2"
}
]
plugin_output = "kudu"
}
}

transform {
}

sink {
Assert {
rules {
table-names = ["kudu_source_table_1", "kudu_source_table_2"]
}
}
}

Table Matching With Regex​

The Kudu Source supports using regular expressions on table_name to match multiple tables (including whole-database style synchronization, since Kudu tables are in a single logical database).

Exact Table Name​

Use table_name to specify a single Kudu table with an exact name:

source {
kudu {
kudu_masters = "kudu-master:7051"
table_name = "kudu_source_table_1"
}
}

Regex Matching​

Use table_name as a regex pattern and enable use_regex to read multiple tables with one configuration:

source {
kudu {
kudu_masters = "kudu-master:7051"
# Match tables like kudu_source_table_1, kudu_source_table_2, etc.
table_name = "kudu_source_table_\\d+"
use_regex = true
}
}

You can also combine regex entries in table_list:

source {
kudu {
kudu_masters = "kudu-master:7051"
table_list = [
{
table_name = "kudu_source_table_1"
},
{
table_name = "kudu_source_table_2"
},
{
# Regex matching - any table whose name starts with prefix_ and ends with digits
table_name = "prefix_\\d+"
use_regex = true
}
]
}
}

Whole-Database Matching​

You can also synchronize all tables in the current Kudu cluster (or all business tables in the current instance, if there are no system tables) by using a catch-all regex:

source {
kudu {
kudu_masters = "kudu-master:7051"
# Match all tables in the current Kudu cluster
table_name = ".*"
use_regex = true
}
}

Changelog​

Change Log
ChangeCommitVersion
[Feature][shade]Refactor the seatunnel-shade module. (#9993)https://github.com/apache/seatunnel/commit/4ba2895953.0.0
[Fix][Connector-V2] Wait for Kudu client shutdown (#11585)https://github.com/apache/seatunnel/commit/4c4b818593.0.0
[Improve][Connector-V2] Improve source split round-robin assignment for InfluxDB IoTDB Kudu and TDengine (#11573)https://github.com/apache/seatunnel/commit/1fe15cbb83.0.0
[Improve][Connector-V2] Upgrade kudu-client from 1.11.1 to 1.15.0 (#10974)https://github.com/apache/seatunnel/commit/477f57f743.0.0
[Chore] fix typos filed -> field (#9757)https://github.com/apache/seatunnel/commit/e3e1c67d292.3.12
[Improve][Core] Update apache common to apache common lang3 (#9694)https://github.com/apache/seatunnel/commit/6e5737c1ec2.3.12
[Improve][API] Optimize the enumerator API semantics and reduce lock calls at the connector level (#9671)https://github.com/apache/seatunnel/commit/9212a771402.3.12
[Feature][connector-kudu] implement the filter (#9405)https://github.com/apache/seatunnel/commit/2714dd11052.3.12
[Feature][Checkpoint] Add check script for source/sink state class serialVersionUID missing (#9118)https://github.com/apache/seatunnel/commit/4f5adeb1c72.3.11
[Improve] kudu options (#9162)https://github.com/apache/seatunnel/commit/e7edafdbac2.3.11
[Improve] restruct connector common options (#8634)https://github.com/apache/seatunnel/commit/f3499a6eeb2.3.10
[Improve][Transform] Rename sql transform table name from 'fake' to 'dual' (#8298)https://github.com/apache/seatunnel/commit/e6169684fb2.3.9
[Improve][dist]add shade check rule (#8136)https://github.com/apache/seatunnel/commit/51ef8000162.3.9
[Improve][API] Unified tables_configs and table_list (#8100)https://github.com/apache/seatunnel/commit/84c0b8d6602.3.9
[Feature][Core] Rename result_table_name/source_table_name to plugin_input/plugin_output (#8072)https://github.com/apache/seatunnel/commit/c7bbd322db2.3.9
[Feature][Restapi] Allow metrics information to be associated to logical plan nodes (#7786)https://github.com/apache/seatunnel/commit/6b7c53d03c2.3.9
[Improve][Connector] Add multi-table sink option check (#7360)https://github.com/apache/seatunnel/commit/2489f6446b2.3.7
[Feature][Core] Support using upstream table placeholders in sink options and auto replacement (#7131)https://github.com/apache/seatunnel/commit/c4ca74122c2.3.6
correct the typo of kudu kerberos config (#6905)https://github.com/apache/seatunnel/commit/fcb85549722.3.6
[Fix][KuduCatalogFactory]: Fix KuduCatalogFactory.optionRule() will throw an Exception (#6787)https://github.com/apache/seatunnel/commit/45a4e1532d2.3.6
[Feature][Engine] Unify job env parameters (#6003)https://github.com/apache/seatunnel/commit/2410ab38f02.3.4
[Feature][Connector-V2] Support multi-table sink feature for kudu (#5951)https://github.com/apache/seatunnel/commit/82460c0bf02.3.4
[Feature] Add unsupported datatype check for all catalog (#5890)https://github.com/apache/seatunnel/commit/b9791285a02.3.4
[Feature][Kudu] Support multi-table source read (#5878)https://github.com/apache/seatunnel/commit/8d9a0b7d112.3.4
[Improve][Common] Introduce new error define rule (#5793)https://github.com/apache/seatunnel/commit/9d1b2582b22.3.4
[Feature][Connector-V2] Support TableSourceFactory/TableSinkFactory on kudu (#5789)https://github.com/apache/seatunnel/commit/10e791d60a2.3.4
[Improve] Remove use SeaTunnelSink::getConsumedType method and mark it as deprecated (#5755)https://github.com/apache/seatunnel/commit/8de74081002.3.4
[Feature][Kudu] Refactor Kudu functionality and Sink support CDC data. (#5437)https://github.com/apache/seatunnel/commit/22110eb7b32.3.4
[Improve][build] Give the maven module a human readable name (#4114)https://github.com/apache/seatunnel/commit/d7cd6010512.3.1
[Improve][Project] Code format with spotless plugin. (#4101)https://github.com/apache/seatunnel/commit/a2ab1665612.3.1
[Hotfix][Connector-V2] Fix connector source snapshot state NPE (#4027)https://github.com/apache/seatunnel/commit/e39c4988cc2.3.1
[Feature][Connector] add get source method to all source connector (#3846)https://github.com/apache/seatunnel/commit/417178fb842.3.1
[Feature][API &amp; Connector &amp; Doc] add parallelism and column projection interface (#3829)https://github.com/apache/seatunnel/commit/b9164b8ba12.3.1
[Hotfix][OptionRule] Fix option rule about all connectors (#3592)https://github.com/apache/seatunnel/commit/226dc6a1192.3.0
[Improve][Connector-V2] Bad smell ToArrayCallWithZeroLengthArrayArgument: (#3577)https://github.com/apache/seatunnel/commit/cc448d98c42.3.0
[Improve][Connector-V2][Kudu] Unified exception for kudu source & sink connector (#3564)https://github.com/apache/seatunnel/commit/273418ddc92.3.0
[Connector][Dependency] Add Miss Dependency Cassandra And Change Kudu Plugin Name (#3432)https://github.com/apache/seatunnel/commit/6ac6a0a0cd2.3.0
[Feature][Connector V2] expose configurable options in Kudu (#3365)https://github.com/apache/seatunnel/commit/c422210e2c2.3.0
[Feature][Core][Connector-V2] Unified The way of setting JobName (#2908)https://github.com/apache/seatunnel/commit/bf2c97484b2.3.0-beta
remove duplicate ExceptionUtil class (#3037)https://github.com/apache/seatunnel/commit/c9dc7c50c22.3.0-beta
[Improve][all] change Log to @Slf4j (#3001)https://github.com/apache/seatunnel/commit/6016100f122.3.0-beta
[Improve][Connector-V2]Kudu Sink Connector Support to upsert rowhttps://github.com/apache/seatunnel/commit/1ece805ab12.3.0-beta
[DEV][Api] Replace SeaTunnelContext with JobContext and remove singleton pattern (#2706)https://github.com/apache/seatunnel/commit/cbf82f755c2.2.0-beta
[#2606]Dependency management split (#2630)https://github.com/apache/seatunnel/commit/fc047be69b2.2.0-beta
[Connector-V2] Add Kudu source and sink connector (#2254)https://github.com/apache/seatunnel/commit/0483cbc2df2.2.0-beta