如何将Kafka Topic中的TIMESTAMP字段转换为ISO 8601格式?
解决方案
方法1:在KSQL中转换字段后输出到新Topic
如果已经创建了基于源Topic的STREAM或TABLE,可通过KSQL内置函数将TIMESTAMP类型(或数字时间戳)转换为ISO 8601格式的字符串,生成新的STREAM/TABLE并输出到目标Topic,供Elasticsearch Sink连接器使用。
假设源TABLE名为USER_TABLE,包含created_at字段(TIMESTAMP类型或数字时间戳),执行以下SQL创建转换后的TABLE:
CREATE TABLE USER_TABLE_ISO WITH ( KAFKA_TOPIC='table-api_auth-user-iso', VALUE_FORMAT='JSON' ) AS SELECT user_id, username, -- 若created_at是TIMESTAMP类型,直接用TO_STRING指定格式 TO_STRING(created_at, 'yyyy-MM-dd''T''HH:mm:ss.SSS''Z''') AS created_at, -- 若created_at是毫秒数字时间戳,先转成TIMESTAMP再格式化 -- TO_STRING(TIMESTAMPFROMUNIXTIME(created_at), 'yyyy-MM-dd''T''HH:mm:ss.SSS''Z''') AS created_at email FROM USER_TABLE EMIT CHANGES;
之后将Elasticsearch Sink连接器的topics配置改为table-api_auth-user-iso即可。
方法2:在JDBC源连接器阶段直接转换
在拉取PostgreSQL数据的JDBC源连接器中,使用Kafka Connect的TimestampConverter转换器,直接将TIMESTAMP字段转换为ISO 8601格式的字符串写入源Topic,后续KSQL无需额外处理。
修改JDBC源连接器配置,添加以下转换规则:
# 启用转换器 transforms=convertTimestamp # 指定转换器类型为值转换 transforms.convertTimestamp.type=org.apache.kafka.connect.transforms.TimestampConverter$Value # 要转换的字段名(对应PostgreSQL的TIMESTAMP字段) transforms.convertTimestamp.field=created_at # 指定ISO 8601格式 transforms.convertTimestamp.format=yyyy-MM-dd'T'HH:mm:ss.SSS'Z' # 转换后的目标类型为字符串 transforms.convertTimestamp.target.type=string
重新启动JDBC源连接器后,源Topictable-api_auth-user中的created_at字段将直接以ISO 8601字符串存储。
方法3:创建KSQL STREAM/TABLE时直接转换字段
如果是从源Topic创建KSQL STREAM/TABLE的阶段,可直接在定义时完成转换:
CREATE STREAM USER_STREAM_ISO WITH ( KAFKA_TOPIC='table-api_auth-user-iso', VALUE_FORMAT='JSON' ) AS SELECT user_id, username, -- 根据源字段类型选择转换方式 TIMESTAMPTOSTRING(created_at, 'yyyy-MM-dd''T''HH:mm:ss.SSS''Z''') AS created_at, email FROM USER_STREAM EMIT CHANGES;
注:TIMESTAMPTOSTRING函数可直接将TIMESTAMP类型或数字时间戳转换为指定格式的字符串。
内容的提问来源于stack exchange,提问作者Leandro Santiago Gomes
相关产品推荐
相关产品推荐

