InfluxDB
InfluxDB source connector
Support Those Engines
Spark
Flink
SeaTunnel Zeta
Description
Read data from InfluxDB 1.x by using an InfluxQL query. The connector supports a normal single
query and an optional parallel scan mode that splits one query by an integer column range.
Use tables_configs to read multiple queries, databases, and output schemas through one source.
Key Features
Data Type Mapping
| SeaTunnel Data Type | Notes |
|---|---|
| BOOLEAN | Parsed from the returned InfluxDB value. |
| SMALLINT | Parsed from the returned InfluxDB value. |
| INT | Parsed from the returned InfluxDB value. |
| BIGINT | Parsed from the returned InfluxDB value. |
| FLOAT | InfluxDB returns numbers as double values; the connector converts them to FLOAT. |
| DOUBLE | Uses the returned numeric value. |
| STRING | Uses the returned value as a string. |
Other SeaTunnel types are not supported by the current InfluxDB source converter.
Source Options
| name | type | required | default value | description |
|---|---|---|---|---|
| url | string | yes | - | InfluxDB server URL, for example http://influxdb-host:8086. |
| sql | string | no | - | Required in single-table mode; mutually exclusive with tables_configs. |
| schema | config | no | - | Required in single-table mode; define it per entry in multi-table mode. |
| database | string | no | - | Required in single-table mode; optional default for multi-table entries. |
| tables_configs | list | no | - | Queries and schemas to read in multi-table mode; see below. |
| username | string | no | - | InfluxDB username. It must be configured together with password. |
| password | string | no | - | InfluxDB password. It must be configured together with username. |
| lower_bound | int | no | - | Lower bound of split_column when parallel scan is enabled. |
| upper_bound | int | no | - | Upper bound of split_column when parallel scan is enabled. |
| partition_num | int | no | 0 | Number of query splits. 0 means the source runs the original sql as one split. |
| split_column | string | no | - | Integer column used to split the query when parallel scan is enabled. |
| where | string | no | - | Reserved source option. The current split logic reads the lowercase where keyword from sql directly. |
| epoch | string | no | n | Time precision returned by InfluxDB. For example: H, m, s, MS, u, n. |
| connect_timeout_ms | long | no | 15000 | Timeout for connecting to InfluxDB, in milliseconds. |
| query_timeout_sec | int | no | 3 | Timeout for querying InfluxDB, in seconds. |
| common-options | config | no | - | Source plugin common options. See Source Common Options. |
url [string]
The URL to connect to InfluxDB, for example http://influxdb-host:8086.
tables_configs [list]
An alternative to the root-level sql and schema. Each entry requires:
sql: its InfluxQL query.database: its database, or inherit the root-leveldatabase.schema: its output fields and a non-blankschema.tablethat uniquely identifies the output table.
An entry can also configure split_column, lower_bound, upper_bound, and partition_num
together. Configure range options inside each entry, not at root level. In this mode the inclusive
integer range is divided without overlaps, including uneven ranges; no more than one split per
integer value is created. The existing single-table range behavior is unchanged.
Partitioned entries support simple SELECT fields FROM measurement queries with an optional
WHERE predicate (case-insensitive). Quoted identifiers, string literals, functions, subqueries,
multiple statements, and trailing clauses such as LIMIT, ORDER BY, or tz(...) are rejected
for range splitting. Use an unpartitioned entry for these queries; it is sent unchanged.
All entries share the root connection options: url, credentials, epoch, and timeouts.
Connection options inside entries are rejected. Root-level sql or schema, an empty list,
missing table names, and duplicate output table identities are rejected before connecting.
The table identity is independent of the measurement name in the query and is used for downstream
table routing. Keep these identities and their query/schema definitions stable when restoring a
checkpoint; changing from single-table to multi-table mode requires a fresh job.
Each query must return the columns declared in its own schema. Empty query results are allowed.
Queries without range splitting, including tz(...), are sent unchanged to InfluxDB.
sql [string]
The InfluxQL query used to read data. For example:
select name, age from test
schema [config]
The output schema of the source. See Schema Feature for the full grammar. For example:
schema {
fields {
name = string
age = int
}
}
database [string]
The InfluxDB database name.
username [string]
InfluxDB username used to authenticate the connection. Configure it together with password.
password [string]
InfluxDB password used to authenticate the connection. Configure it together with username.
split_column [string]
The integer column used to split one query into multiple range queries when parallel scan is enabled.
Tips:
- InfluxDB tags cannot be used as a split column because tags only support the string type.
- InfluxDB time cannot be used as a split column because the time field cannot participate in mathematical calculations.
split_columncurrently only supports integer columns;float,string,date, and other types are not supported.split_column,lower_bound,upper_bound, andpartition_nummust be configured together.- If the split query has a filter, write the filter with a lowercase
wheredirectly insidesql(for exampleselect * from test where age > 0). The current split parser is case-sensitive.whereis part of the option validation rule but the split logic reads the filter fromsql. Put the filter insqlinstead of configuring a separatewherevalue.
upper_bound [int]
Upper bound of the split_column column when parallel scan is enabled.
lower_bound [int]
Lower bound of the split_column column when parallel scan is enabled.
The split column range is divided into partition_num parts. If partition_num = 1, the connector
uses the whole range. If partition_num is less than upper_bound - lower_bound, the connector
uses (upper_bound - lower_bound) partitions.
For example, with lower_bound = 1, upper_bound = 10, partition_num = 2, and
sql = "select * from test where age > 0 and age < 10", the connector splits the query into:
split 1: select * from test where ($split_column >= 1 and $split_column < 6) and ( age > 0 and age < 10 )
split 2: select * from test where ($split_column >= 6 and $split_column < 11) and ( age > 0 and age < 10 )
partition_num [int]
Number of query splits. Configure it together with lower_bound, upper_bound, and
split_column. Make sure upper_bound - lower_bound is divisible by partition_num; otherwise
the query results overlap.
where [string]
The current split logic reads the lowercase where keyword directly from the sql option.
Setting where has no effect on the generated split queries — put the filter inside sql
instead, for example select * from test where age > 0. The split parser is case-sensitive on
the where keyword.
This option is kept in the option list only because the validation rule
(InfluxDBSourceFactory.optionRule()) still references it for backward compatibility; the
runtime split query generator does not read it.
epoch [string]
Time precision returned by InfluxDB. Valid values include H, m, s, MS, u, and n. The
default value is n.
query_timeout_sec [int]
Query timeout for the InfluxDB client, in seconds.
connect_timeout_ms [long]
Connection timeout for the InfluxDB client, in milliseconds.
common options
Source plugin common parameters, please refer to Source Common Options for details.
Task Example
Read Multiple Tables
env {
parallelism = 2
job.mode = "BATCH"
}
source {
InfluxDB {
url = "http://influxdb-host:8086"
tables_configs = [
{
database = "telemetry"
sql = "select value from temperature"
schema {
table = "temperatures"
fields { value = DOUBLE }
}
},
{
database = "operations"
sql = "select active, label from alerts"
schema {
table = "alerts"
fields {
active = BOOLEAN
label = STRING
}
}
}
]
}
}
sink { Console {} }
Read With Parallel Range Splits
env {
parallelism = 1
job.mode = "BATCH"
}
source {
InfluxDB {
url = "http://influxdb-host:8086"
sql = "select label, c_string, c_double, c_bigint, c_float, c_int, c_smallint, c_boolean from source"
database = "test"
upper_bound = 99
lower_bound = 0
partition_num = 4
split_column = "c_int"
schema {
fields {
label = STRING
c_string = STRING
c_double = DOUBLE
c_bigint = BIGINT
c_float = FLOAT
c_int = INT
c_smallint = SMALLINT
c_boolean = BOOLEAN
time = BIGINT
}
}
}
}
sink {
Console {}
}
Read Without Parallel Range Splits
env {
parallelism = 1
job.mode = "BATCH"
}
source {
InfluxDB {
url = "http://influxdb-host:8086"
sql = "select label, c_string, c_double, c_bigint, c_float, c_int, c_smallint, c_boolean from source"
database = "test"
schema {
fields {
label = STRING
c_string = STRING
c_double = DOUBLE
c_bigint = BIGINT
c_float = FLOAT
c_int = INT
c_smallint = SMALLINT
c_boolean = BOOLEAN
time = BIGINT
}
}
}
}
sink {
Console {}
}
Read With InfluxQL Time Zone
env {
parallelism = 1
job.mode = "BATCH"
}
source {
InfluxDB {
url = "http://influxdb-host:8086"
sql = "select label, c_string, c_double, c_bigint, c_float, c_int, c_smallint, c_boolean from source tz('Asia/Shanghai')"
database = "test"
schema {
fields {
label = STRING
c_string = STRING
c_double = DOUBLE
c_bigint = BIGINT
c_float = FLOAT
c_int = INT
c_smallint = SMALLINT
c_boolean = BOOLEAN
time = BIGINT
}
}
}
}
sink {
Console {}
}
Bounded Read With Parallel Splits
The InfluxDB source is bounded by the sql result. To chain it to a streaming sink without
losing rows, set a finite partition_num and let SeaTunnel checkpoint the offsets per split.
env {
parallelism = 2
job.mode = "BATCH"
checkpoint.interval = 10000
}
source {
InfluxDB {
url = "http://influxdb-host:8086"
sql = "select label, c_string, c_double, c_bigint, c_float, c_int, c_smallint, c_boolean from source"
database = "test"
lower_bound = 0
upper_bound = 99
partition_num = 4
split_column = "c_int"
query_timeout_sec = 10
connect_timeout_ms = 20000
schema {
fields {
label = STRING
c_string = STRING
c_double = DOUBLE
c_bigint = BIGINT
c_float = FLOAT
c_int = INT
c_smallint = SMALLINT
c_boolean = BOOLEAN
time = BIGINT
}
}
}
}
sink {
InfluxDB {
url = "http://influxdb-host:8086"
database = "test"
measurement = "sink"
key_time = "time"
key_tags = ["label"]
batch_size = 1024
}
}
Changelog
Change Log
| Change | Commit | Version |
|---|---|---|
| [Feature][Checkpoint] Add check script for source/sink state class serialVersionUID missing (#9118) | https://github.com/apache/seatunnel/commit/4f5adeb1c7 | 2.3.11 |
| [Improve] influxdb options (#8966) | https://github.com/apache/seatunnel/commit/9f498b8133 | 2.3.10 |
| [Improve] restruct connector common options (#8634) | https://github.com/apache/seatunnel/commit/f3499a6eeb | 2.3.10 |
| [Improve][dist]add shade check rule (#8136) | https://github.com/apache/seatunnel/commit/51ef800016 | 2.3.9 |
| [Feature][Restapi] Allow metrics information to be associated to logical plan nodes (#7786) | https://github.com/apache/seatunnel/commit/6b7c53d03c | 2.3.9 |
| [Improve] Improve some connectors prepare check error message (#7465) | https://github.com/apache/seatunnel/commit/6930a25edd | 2.3.8 |
| [Improve][Connector] Add multi-table sink option check (#7360) | https://github.com/apache/seatunnel/commit/2489f6446b | 2.3.7 |
| [Feature][Core] Support using upstream table placeholders in sink options and auto replacement (#7131) | https://github.com/apache/seatunnel/commit/c4ca74122c | 2.3.6 |
| Support multi-table sink feature for influxdb (#6278) | https://github.com/apache/seatunnel/commit/56f13e920d | 2.3.5 |
| [Improve][Zeta] Add classloader cache mode to fix metaspace leak (#6355) | https://github.com/apache/seatunnel/commit/9c3c2f183d | 2.3.5 |
| [Test][E2E] Add thread leak check for connector (#5773) | https://github.com/apache/seatunnel/commit/1f2f3fc5f0 | 2.3.4 |
| [BugFix][InfluxDBSource] Resolve invalid SQL in initColumnsIndex method caused by direct QUERY_LIMIT appendage with 'tz' function. (#4829) | https://github.com/apache/seatunnel/commit/deed9c62c3 | 2.3.4 |
| [Improve][Common] Introduce new error define rule (#5793) | https://github.com/apache/seatunnel/commit/9d1b2582b2 | 2.3.4 |
[Improve] Remove use SeaTunnelSink::getConsumedType method and mark it as deprecated (#5755) | https://github.com/apache/seatunnel/commit/8de7408100 | 2.3.4 |
| Support config column/primaryKey/constraintKey in schema (#5564) | https://github.com/apache/seatunnel/commit/eac76b4e50 | 2.3.4 |
| [Improve][Connector-V2] Remove scheduler in InfluxDB sink (#5271) | https://github.com/apache/seatunnel/commit/f459f500cb | 2.3.4 |
| [Improve][CheckStyle] Remove useless 'SuppressWarnings' annotation of checkstyle. (#5260) | https://github.com/apache/seatunnel/commit/51c0d709ba | 2.3.4 |
| Merge branch 'dev' into merge/cdc | https://github.com/apache/seatunnel/commit/4324ee1912 | 2.3.1 |
| [Improve][Project] Code format with spotless plugin. | https://github.com/apache/seatunnel/commit/423b583038 | 2.3.1 |
| [improve][api] Refactoring schema parse (#4157) | https://github.com/apache/seatunnel/commit/b2f573a13e | 2.3.1 |
| [Improve][build] Give the maven module a human readable name (#4114) | https://github.com/apache/seatunnel/commit/d7cd601051 | 2.3.1 |
| [Improve][Project] Code format with spotless plugin. (#4101) | https://github.com/apache/seatunnel/commit/a2ab166561 | 2.3.1 |
| [Improve][SourceConnector] Unifie InfluxDB source fields to schema (#3897) | https://github.com/apache/seatunnel/commit/85a984a64f | 2.3.1 |
| [Feature][Connector] add get source method to all source connector (#3846) | https://github.com/apache/seatunnel/commit/417178fb84 | 2.3.1 |
| [Feature][API & Connector & Doc] add parallelism and column projection interface (#3829) | https://github.com/apache/seatunnel/commit/b9164b8ba1 | 2.3.1 |
| [Hotfix][OptionRule] Fix option rule about all connectors (#3592) | https://github.com/apache/seatunnel/commit/226dc6a119 | 2.3.0 |
| [Improve][Connector-V2][Influxdb] Unified exception for influxdb source & sink connector (#3558) | https://github.com/apache/seatunnel/commit/4686f35d68 | 2.3.0 |
| [Feature][Connector][influx] Expose configurable options in influx db (#3392) | https://github.com/apache/seatunnel/commit/b247ff0aef | 2.3.0 |
| [Feature][Connector-V2] influxdb sink connector (#3174) | https://github.com/apache/seatunnel/commit/630e884791 | 2.3.0 |
| [Feature][Connector-V2] Add influxDB connector source (#2697) | https://github.com/apache/seatunnel/commit/1d70ea3084 | 2.3.0-beta |