帮你快速理解、总结文档立即下载
文档中心>流计算 Oceanus>SQL 开发指南>上下游开发指南>数据分析引擎 Elasticsearch Service

数据分析引擎 Elasticsearch Service

最近更新时间:2026-08-07 17:16:01
我的收藏

介绍

Elasticsearch Connector 提供了对 Elasticsearch 的写入、读取和维表关联支持。目前 Oceanus 支持 Elasticsearch 6.x、7.x 和 8.x 版本。

版本说明

Flink 版本
说明
ES 版本
1.13
写入、批数据源
6.x、7.x
1.14
写入
6.x、7.x
1.16
写入
6.x、7.x
1.18
写入、维表
6.x、7.x、8.x
1.20
写入、维表
6.x、7.x、8.x
注意:
Flink 1.18 及以上版本中,Elasticsearch Connector 基于新的 SinkV2 API 实现。ES6/ES7 的参数与旧版本兼容,ES8 使用独立的 connector 标识 elasticsearch-8,参数体系有较大变化,详见下方各节说明。

使用范围

Elasticsearch 支持写入,可以作为 Tuple 数据流的目的表(Sink),也可以作为 Upsert 数据流的目的表(Sink,自动以文档 _id 字段生成主键,并更新之前的文档版本)。
如果希望将 JDBC 数据库的变动记录,将其作为流式源表消费,可以使用 DebeziumCanal 等,对 JDBC 数据库的变更进行捕获和订阅,然后 Flink 即可对这些变更事件进行进一步的处理。可参见 Kafka
Flink 1.13 支持 Elasticsearch 的批模式读,仅支持 Elasticsearch 7。
Flink 1.18 及以上版本支持将 Elasticsearch 作为维表(Lookup Table),用于流表 JOIN 关联查询,支持 ES 6.x、7.x 和 8.x 版本。

DDL 定义

用作 Elasticsearch 6 数据目的(Sink)

CREATE TABLE elasticsearch6_sink_table (
`id` INT,
`name` STRING,
PRIMARY KEY (`id`) NOT ENFORCED -- 对应 Elasticsearch 中的 _id
) WITH (
'connector' = 'elasticsearch-6', -- 输出到 Elasticsearch 6
'username' = '$username', -- 选填 用户名
'password' = '$password', -- 选填 密码
'hosts' = 'http://10.28.28.94:9200', -- Elasticsearch 的连接地址
'index' = 'my-index', -- Elasticsearch 的 Index 名
'document-type' = '_doc', -- Elasticsearch 的 Document 类型
'format' = 'json' -- 输出数据格式,目前只支持 'json'
);

用作 Elasticsearch 7 数据目的(Sink)

CREATE TABLE elasticsearch7_sink_table (
`id` INT,
`name` STRING,
PRIMARY KEY (`id`) NOT ENFORCED -- 对应 Elasticsearch 中的 _id
) WITH (
'connector' = 'elasticsearch-7', -- 输出到 Elasticsearch 7
'username' = '$username', -- 选填 用户名
'password' = '$password', -- 选填 密码
'hosts' = 'http://10.28.28.94:9200', -- Elasticsearch 的连接地址
'index' = 'my-index', -- Elasticsearch 的 Index 名
'format' = 'json' -- 输出数据格式,目前只支持 'json'
);

作为 Elasticsearch 7 批数据源(Source)

CREATE TABLE elasticsearch7_source_table (
`id` bigint,
`event_date` int,
`app` int,
primary key (`id`) not enforced
) with (
-- 必填参数
'connector' = 'es-source',
'endPoint' = '127.0.0.1', -- Elasticsearch 的连接 ip
'accessId' = 'elastic', -- 用户名
'accessKey' = 'PASSWORD', -- 密码
'indexName' = 'my-index', -- Elasticsearch 的 Index 名
'format' = 'json', -- 数据格式,只支持 'json'
-- 可选参数
'scheme' = 'http', -- 连接协议
'port' = '9200', -- 端口
'batchSize' = '2000', -- 每个 scroll 请求从 Elasticsearch 集群获取的最大文档数
'keepScrollAliveSecs' = '60' -- scroll 上下文保留的最长时间,单位为分钟
);

用作 Elasticsearch 8 数据目的(Sink)

说明:
适用于 Flink 1.18 及以上版本。
CREATE TABLE elasticsearch8_sink_table (
`id` INT,
`name` STRING,
PRIMARY KEY (`id`) NOT ENFORCED -- 对应 Elasticsearch 中的 _id
) WITH (
'connector' = 'elasticsearch-8', -- 输出到 Elasticsearch 8
'username' = '$username', -- 选填 用户名
'password' = '$password', -- 选填 密码
'hosts' = 'http://10.28.28.94:9200', -- Elasticsearch 的连接地址
'index' = 'my-index', -- Elasticsearch 的 Index 名
'format' = 'json' -- 输出数据格式,目前只支持 'json'
);

用作 Elasticsearch 8 维表(Lookup Table)

说明:
适用于 Flink 1.18 及以上版本。支持通过维表 JOIN 方式关联查询 Elasticsearch 中的数据进行字段补全。
CREATE TABLE elasticsearch8_lookup_table (
`id` INT,
`name` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'elasticsearch-8', -- ES8 维表
'username' = '$username', -- 选填 用户名
'password' = '$password', -- 选填 密码
'hosts' = 'http://10.28.28.94:9200', -- Elasticsearch 的连接地址
'index' = 'my-index', -- Elasticsearch 的 Index 名
'format' = 'json', -- 输出数据格式,目前只支持 'json'
'lookup.cache' = 'PARTIAL', -- 维表缓存策略,可选 ALL / PARTIAL / NONE
'lookup.partial-cache.max-rows' = '10000', -- PARTIAL 缓存最大行数
'lookup.partial-cache.expire-after-access' = '30min', -- PARTIAL 缓存访问过期时间
'lookup.partial-cache.expire-after-write' = '30min' -- PARTIAL 缓存写入过期时间
);

WITH 参数

作为数据目的(Elasticsearch 6/7)

说明:
此参数表同时适用于 Flink 1.16 及之前版本和 Flink 1.18+ 版本(ES6/ES7 的 elasticsearch-6 / elasticsearch-7 connector),两个版本参数基本一致。
参数值
必填
默认值
描述
connector
当写入 Elasticsearch 6.x 版本时,取值 elasticsearch-6。当写入 Elasticsearch 7.x 及以上版本时,取值 elasticsearch-7
username
用户名。
password
密码。
hosts
Elasticsearch 的连接地址。
index
数据要写入的 Index。支持固定 Index(例如 'myIndex'),也支持动态 Index(例如'index-{log_ts\\|yyyy-MM-dd}')。
document-type
6.x 版本:必填
7.x 版本:不需要
Elasticsearch 文档的 Type 信息。当选择 elasticsearch-7 时,不能填写这个字段,否则会报错。
document-id.key-delimiter
_
为复合键生成 _id 时的分隔符 (默认是 "_")。例如有 a、b、c 三个主键,某条数据的 a 字段为 "1",b 字段为 "2",c 字段为 "3",使用默认分隔符,则最终写入 Elasticsearch 的 _id 是 "1_2_3"。
drop-delete
false
是否过滤上游传来的 DELETE(删除)消息。
此外,在多表 LEFT JOIN 且 JOIN Key 非主键的场景下,启用该选项后,可以解决 Elasticsearch 收到较多临时 null 值数据的问题。需要注意的是,JOIN 左右表的字段不能含有 null 值,否则可能会丢失部分数据。
failure-handler
fail
指定请求 Elasticsearch 失败时,错误处理策略。选项为:
- fail:抛出一个异常。
- ignore:忽略错误,直接继续。
- retry-rejected:重试写入该条记录。
- 另外也支持自定义错误处理器,这里可以填写用户自己编写的 Handler 的类全名(需要上传自定义程序包)。
sink.flush-on-checkpoint
true
Flink 进行快照时,是否等待现有记录完全写入 Elasticsearch 。如果设置为 false,则可能造成恢复时部分数据丢失或者重复等异常情况,但快照速度会提升。
sink.bulk-flush.max-actions
1000
批量写入的最大条数。设置为 0 则禁用批量功能。
sink.bulk-flush.max-size
2MB
批量写入缓存的最大容量,必须以 MB 为单位。设置为 0 则禁用批量功能。
sink.bulk-flush.interval
1s
批量写入的刷新周期。设置为0则禁用批量功能。
sink.bulk-flush.backoff.strategy
DISABLED
批量写入时,失败重试的策略。
- DISABLED:不重试。
- CONSTANT:等待 sink.bulk-flush.backoff.delay 选项设置的毫秒后重试。
- EXPONENTIAL:一开始等待 sink.bulk-flush.backoff.delay 选项设置的毫秒后重试,每次失败后将指数增加下次的等待时间。
sink.bulk-flush.backoff.max-retries
8
批量写入时,最多失败重试的次数。
sink.bulk-flush.backoff.delay
50ms
批量写入失败时,每次重试之间的等待间隔(对于 CONSTANT 策略而言)或间隔的初始基数(对于 EXPONENTIAL 策略而言)。
connection.max-retry-timeout
重试请求的最大超时时间,例如:"20 s"。
connection.path-prefix
指定每个 REST 请求的前缀,例如 '/v1'。通常不需要设置该选项。
format
json
指定输出的格式,默认是内置的 json 格式,可以使用 前文(Kafka)描述过的 JSON 格式选项,例如 json.fail-on-missing-fieldjson.ignore-parse-errorsjson.timestamp-format.standard 等。
retry-on-conflict
更新操作中,允许因版本冲突异常而重试的最大次数。超过该次数后将抛出异常导致作业失败。
routing-key
可以指定分片路由字段,例如 user_id、name 等。
ssl.truststore.path
仅 Flink 1.18+ 支持。HTTPS 连接时 JKS/PKCS12 trust store 文件路径。当连接使用自签名或内部 CA 证书的 ES 集群时需设置此参数。
ssl.truststore.password
仅 Flink 1.18+ 支持。Trust store 文件的密码。
ssl.verification-mode
full
仅 Flink 1.18+ 支持。SSL 验证模式:full(验证证书和主机名)、certificate(仅验证证书)、none(不验证,不推荐生产使用)。

作为数据目的(Elasticsearch 8)

说明:
适用于 Flink 1.18 及以上版本,connector 标识为 elasticsearch-8
参数值
必填
默认值
描述
connector
固定值 elasticsearch-8
hosts
Elasticsearch 的连接地址,格式为 http://host_name:port
index
数据要写入的 Index。支持固定 Index(例如 'myIndex'),也支持动态 Index(例如 'index-{log_ts|yyyy-MM-dd}')。
username
用户名。
password
密码。
document-id.key-delimiter
_
为复合键生成 _id 时的分隔符 (默认是 "_")。例如有 a、b、c 三个主键,某条数据的 a 字段为 "1",b 字段为 "2",c 字段为 "3",使用默认分隔符,则最终写入 Elasticsearch 的 _id 是 "1_2_3"。
format
json
指定输出的格式,默认是内置的 json 格式,可以使用 前文(Kafka)描述过的 JSON 格式选项,例如 json.fail-on-missing-fieldjson.ignore-parse-errorsjson.timestamp-format.standard 等。
sink.delivery-guarantee
AT_LEAST_ONCE
写入的交付语义。可选值:AT_LEAST_ONCEEXACTLY_ONCENONE
sink.bulk-flush.max-actions
1000
每个批量请求的最大操作数。
sink.bulk-flush.max-size
2mb
批量请求中缓冲操作的最大大小,必须以 MB 为单位。
sink.bulk-flush.interval
1s
批量 flush 的时间间隔。
sink.bulk-flush.max-buffered-actions
10000
缓冲区中最大缓冲的操作数。超过该值后,新的写入将被阻塞直到缓冲区被消费。
sink.bulk-flush.max-in-flight-actions
50
同时处于飞行中(未完成)的最大操作数。
retry-on-conflict
更新操作中,允许因版本冲突异常而重试的最大次数。超过该次数后将抛出异常导致作业失败。
routing-key
可以指定分片路由字段,如 user_id、name 等。Elasticsearch 会自动对 routing-key 进行哈希并将其分配到对应分片。
connection.path-prefix
指定每个 REST 请求的前缀,例如 '/v1'。通常不需要设置该选项。
connection.request-timeout
从连接管理器请求连接的超时时间,例如:"30 s"。
connection.timeout
建立连接的超时时间,例如:"30 s"。
socket.timeout
等待数据的 Socket 超时(SO_TIMEOUT),即两个连续数据包之间的最大不活跃时间,例如:"30 s"。
ssl.certificate-fingerprint
HTTPS 连接的 CA 证书 SHA-256 指纹,用于验证 HTTPS 连接。
ssl.truststore.path
HTTPS 连接的 JKS/PKCS12 trust store 文件路径。当连接使用自签名或内部 CA 证书的 ES 集群时需要此参数。
ssl.truststore.password
Trust store 文件的密码。
ssl.verification-mode
full
SSL 验证模式。
- full:验证证书和主机名(默认,推荐生产使用)。
- certificate:只验证证书,跳过主机名验证。
- none:禁用所有 SSL 验证(不推荐用于生产环境)。
sink.parallelism
指定 Sink 算子的并行度。不设置则由 Flink 框架自动推断。

作为数据源

参数值
必填
默认值
描述
connector
固定值 es-source
endPoint
Elasticsearch 的连接 IP,示例 127.0.0.1
accessId
用户名
accessKey
密码
indexName
要读取的 Index
format
指定读取的格式,只支持内置的 json 格式,可以使用 前文(Kafka)描述过的 JSON 格式选项,例如 json.fail-on-missing-fieldjson.ignore-parse-errorsjson.timestamp-format.standard 等。
scheme
http
Elasticsearch 连接模式,例如 httphttps
port
9200
Elasticsearch 连接端口
batchSize
2000
每个 scroll 请求从 Elasticsearch 集群获取的最大文档数
keepScrollAliveSecs
60
scroll 上下文保留的最长时间,单位为分钟

代码示例

作为数据目的

CREATE TABLE datagen_source_table (
id INT,
name STRING
) WITH (
'connector' = 'datagen',
'rows-per-second'='1' -- 每秒产生的数据条数
);

CREATE TABLE elasticsearch7_sink_table (
`id` INT,
`name` STRING
) WITH (
'connector' = 'elasticsearch-7', -- 输出到 Elasticsearch 7
'username' = '$username', -- 选填 用户名
'password' = '$password', -- 选填 密码
'hosts' = 'http://10.28.28.94:9200', -- Elasticsearch 的连接地址
'index' = 'my-index', -- Elasticsearch 的 Index 名
'sink.bulk-flush.max-actions' = '1000', -- 数据刷新频率
'sink.bulk-flush.interval' = '1s' -- 数据刷新周期
'format' = 'json' -- 输出数据格式,目前只支持 'json'
);

INSERT INTO elasticsearch7_sink_table select * from datagen_source_table;

作为数据目的(Elasticsearch 8)

说明:
适用于 Flink 1.18 及以上版本。
CREATE TABLE datagen_source_table (
id INT,
name STRING
) WITH (
'connector' = 'datagen',
'rows-per-second'='1' -- 每秒产生的数据条数
);

CREATE TABLE elasticsearch8_sink_table (
`id` INT,
`name` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'elasticsearch-8', -- 输出到 Elasticsearch 8
'username' = '$username', -- 选填 用户名
'password' = '$password', -- 选填 密码
'hosts' = 'http://10.28.28.94:9200', -- Elasticsearch 的连接地址
'index' = 'my-index', -- Elasticsearch 的 Index 名
'format' = 'json', -- 输出数据格式,目前只支持 'json'
'sink.bulk-flush.max-actions' = '1000', -- 每个批量请求的最大操作数
'sink.bulk-flush.interval' = '1s' -- 批量 flush 的时间间隔
);

INSERT INTO elasticsearch8_sink_table SELECT * FROM datagen_source_table;

作为维表关联查询(Elasticsearch 8 Lookup Join)

说明:
适用于 Flink 1.18 及以上版本。将 Elasticsearch 作为维表,与流表进行 JOIN 查询。
-- 维表:Elasticsearch 中的字典表
CREATE TABLE dim_es_user (
`id` INT,
`name` STRING,
`age` INT,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'elasticsearch-8', -- ES8 维表
'hosts' = 'http://10.28.28.94:9200',
'index' = 'user-dim',
'format' = 'json',
'lookup.cache' = 'PARTIAL', -- 缓存策略:ALL / PARTIAL / NONE
'lookup.partial-cache.max-rows' = '10000', -- PARTIAL 缓存最大行数
'lookup.partial-cache.expire-after-write' = '30min'
);

-- 流表:业务数据流
CREATE TABLE kafka_order (
`order_id` INT,
`user_id` INT,
`amount` DOUBLE,
`ts` TIMESTAMP(3)
) WITH (
'connector' = 'kafka',
-- ... Kafka 连接配置 ...
);

-- 维表 JOIN 查询,用 user_id 关联 ES 维表中的用户信息
INSERT INTO result_sink
SELECT
o.order_id,
o.amount,
u.name,
u.age
FROM kafka_order AS o
LEFT JOIN dim_es_user FOR SYSTEM_TIME AS OF o.ts AS u
ON o.user_id = u.id;

作为批数据源

CREATE TABLE elasticsearch7_source_table (
`id` bigint,
`event_date` int,
`app` int,
primary key (`id`) not enforced
) with (
-- 必填参数
'connector' = 'es-source',
'endPoint' = '127.0.0.1', -- Elasticsearch 的连接 ip
'accessId' = 'elastic', -- 用户名
'accessKey' = 'PASSWORD', -- 密码
'indexName' = 'my-index', -- Elasticsearch 的 Index 名
'format' = 'json', -- 数据格式,只支持 'json'
-- 可选参数
'scheme' = 'http', -- 连接协议
'port' = '9200', -- 端口
'batchSize' = '2000', -- 每个 scroll 请求从 Elasticsearch 集群获取的最大文档数
'keepScrollAliveSecs' = '60' -- scroll 上下文保留的最长时间,单位为分钟
);

CREATE TABLE logger_sink (
`id` bigint,
`event_date` int,
`app` int,
primary key (`id`) not enforced
) with (
'connector' = 'logger'
);

INSERT INTO logger_sink SELECT * from elasticsearch7_source_table;

注意事项

如果您希望连接其他版本的 Elasticsearch,请通过附加自定义程序包的方式,上传相应的 Elasticsearch Sink 的 JAR 包。

监控指标说明

Oceanus 为 ES Connector 增加了很多实用的统计指标。单击 Flink UI 的运行图中的 ES Sink 算子,即可搜索并查看指标:
numberOfInsertRecords:获取输出的 +I 消息数。
numberOfDeleteRecords:获取输出的 -D 消息数。
numberOfUpdateBeforeRecords:获取输出的 -U 消息数。
numberOfUpdateAfterRecords:获取输出的 +U 消息数。