Skip to main content
Version: Next

Enterprise WeChat

Enterprise WeChat sink connector

Support Those Engines

Spark
Flink
SeaTunnel Zeta

Description

A sink plugin that sends SeaTunnel rows to an Enterprise WeChat robot webhook. The connector identifier in job configuration is WeChat. Each row is serialized into a plain-text message with the fields rendered as fieldName: fieldValue lines before the request is sent to the webhook.

The connector extends the Http sink so it inherits the standard HTTP retry behaviour (retry, retry_backoff_multiplier_ms, retry_backoff_max_ms) and the generic multi_table_sink_replica option for fan-out writes.

For example, if the data from upstream is {"alarmStatus": "firing", "alarmTime": "2022-08-03 01:38:49", "alarmContent": "The disk usage exceeds the threshold"}, the output content to WeChat Robot is the following:

alarmStatus: firing
alarmTime: 2022-08-03 01:38:49
alarmContent: The disk usage exceeds the threshold

Key features

Data Type Mapping

The connector renders every upstream row as one plain-text message. Each field is converted to its string representation (String.valueOf(value)) and emitted on its own line together with the field name, so no per-type JSON structure exists on the wire.

SeaTunnel Data TypeEnterprise WeChat Message Field
stringfieldName: string
tinyint / smallint / int / bigintfieldName: number
float / doublefieldName: number
booleanfieldName: true/false
date / time / timestampfieldName: ISO string
bytes / array / map / rowfieldName: String(toString)

Options

nametyperequireddefault valuedescription
urlStringYes-Enterprise WeChat robot webhook URL, format https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=XXXXXX.
mentioned_listarrayNo-User IDs to mention in the group. Use @all to mention everyone.
mentioned_mobile_listarrayNo-Mobile phone numbers to mention in the group. Use @all to mention everyone.
retryintNo-Maximum retry times when the HTTP request throws IOException.
retry_backoff_multiplier_msintNo100Retry backoff multiplier in milliseconds.
retry_backoff_max_msintNo10000Maximum retry backoff in milliseconds.
multi_table_sink_replicaintNo1Number of writer replicas used when writing multiple tables.
common-optionsno-Sink plugin common parameters, please refer to Sink Common Options for details.

url [string]

Enterprise WeChat webhook URL format is https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=XXXXXX. The key query parameter is the robot key generated in the Enterprise WeChat group robot settings.

mentioned_list [array]

A list of user IDs to mention in the group. Use @all to mention everyone. If the user ID is unavailable, use mentioned_mobile_list.

mentioned_mobile_list [array]

Mobile phone numbers to mention in the group. Use @all to mention everyone.

retry [int]

Maximum retry times when the HTTP request throws IOException. There is no retry by default. The retry loop uses retry_backoff_multiplier_ms and retry_backoff_max_ms to compute the wait between attempts.

retry_backoff_multiplier_ms [int]

Base unit (in milliseconds) for the retry backoff. The wait between attempts grows across retries up to retry_backoff_max_ms. The growth curve is not a fixed multiplier per attempt — see HttpClientProvider (connector-http-base) for the exact Fibonacci-based strategy. Default is 100.

retry_backoff_max_ms [int]

Maximum wait between retries, in milliseconds. Default is 10000.

multi_table_sink_replica [int]

Number of writer replicas used when writing multiple tables. Increase this value to add more parallel writers per table. Default is 1.

common options

Sink plugin common parameters, please refer to Sink Common Options for details.

Task Example

Simple

env {
parallelism = 1
job.mode = "BATCH"
}

source {
FakeSource {
row.num = 1
schema = {
fields {
alarmStatus = string
alarmTime = string
alarmContent = string
}
}
rows = [
{
fields = ["firing", "2022-08-03 01:38:49", "The disk usage exceeds the threshold"]
}
]
}
}

sink {
WeChat {
url = "https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=693axxx6-7aoc-4bc4-97a0-0ec2sifa5aaa"
}
}

Mention users and phone numbers

env {
parallelism = 1
job.mode = "BATCH"
}

source {
FakeSource {
row.num = 1
schema = {
fields {
alarmStatus = string
alarmTime = string
alarmContent = string
}
}
rows = [
{
fields = ["firing", "2022-08-03 01:38:49", "The disk usage exceeds the threshold"]
}
]
}
}

sink {
WeChat {
url = "https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=693axxx6-7aoc-4bc4-97a0-0ec2sifa5aaa"
mentioned_list = ["wangqing", "@all"]
mentioned_mobile_list = ["13800001111", "@all"]
}
}

Changelog

Change Log
ChangeCommitVersion
[improve] http connector options (#8969)https://github.com/apache/seatunnel/commit/63ff9f910a2.3.10
[Feature][Restapi] Allow metrics information to be associated to logical plan nodes (#7786)https://github.com/apache/seatunnel/commit/6b7c53d03c2.3.9
[Feature][Core] Support using upstream table placeholders in sink options and auto replacement (#7131)https://github.com/apache/seatunnel/commit/c4ca74122c2.3.6
[Feature][Connector-V2] Support TableSourceFactory/TableSinkFactory on http (#5816)https://github.com/apache/seatunnel/commit/6f49ec6ead2.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][OptionRule] Fix option rule about all connectors (#3592)https://github.com/apache/seatunnel/commit/226dc6a1192.3.0
[Feature][Connector-V2][Http] Add option rules && Improve Myhours sink connector (#3351)https://github.com/apache/seatunnel/commit/cc8bb60c832.3.0
[Bug][Connector-V2] Fix wechat sink data serialization (#2856)https://github.com/apache/seatunnel/commit/3aee11fc162.3.0-beta
[Feature][Connector-V2] Add Enterprise Wechat sink connector (#2412)https://github.com/apache/seatunnel/commit/3e200e0a382.2.0-beta