You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何利用Redis Sink Connector插入Hash类型数据?SMT配置可行性咨询

实现Redis Sink Connector写入Hash类型数据的方案

完全可以通过配置SMT(Single Message Transform)实现需求,核心是将嵌套的JSON结构转换为Redis Hash能识别的扁平键值对,并配合Connector的特定配置完成写入。以下是具体步骤:

1. 用SMT提取嵌套的有效负载

你的示例数据外层被ADDRESSBOOKLIST字段包裹,首先需要提取内层的对象作为消息的value,这样Redis Sink才能将其解析为Hash的键值对。添加以下SMT配置:

transforms=extractPayload
transforms.extractPayload.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.extractPayload.field=ADDRESSBOOKLIST

2. 配置Redis Sink Connector为Hash写入模式

修改Connector的核心配置,指定写入命令为HSET(对应Redis Hash类型),并确保value转换器能正确解析Map结构:

# 指定Redis写入命令为HSET
redis.command=HSET
# 设置Redis中Hash的标识key(可根据业务选择,比如用消息里的userId作为Hash key)
# 若需要从消息中动态提取Hash key,可额外添加SMT提取字段作为record key
redis.key=your-hash-key-pattern
# 配置value转换器为JSON转换器,禁用schema以适配无schema的JSON数据
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false

如果需要动态生成Redis Hash的key(比如用消息中的userId),可以再加一个SMT提取字段作为Kafka record的key:

transforms=extractPayload,setRedisKey
transforms.setRedisKey.type=org.apache.kafka.connect.transforms.ExtractField$Key
transforms.setRedisKey.field=userId

3. 验证处理后的数据结构

经过SMT处理后,消息的value会变成扁平的键值对:

{
    "contactId": 123456,
    "userId": 112584,
    "companyName": "Other",
    "firstName": "sam",
    "lastName": "william",
    "contactType": "CARRIER"
}

此时Redis Sink会自动将每个键值对映射为Hash的field和value完成写入。

注意事项

你之前设置的KEY_FORMAT=json是针对Kafka record key的格式配置,和Redis写入的value类型无关,所以需要通过上述SMT和Connector配置调整才能实现Hash类型的写入。

内容的提问来源于stack exchange,提问作者Santhosh S

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.29 03:52:58