ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

SeaTunnel Redis Source 连接器实战指南:从单表 SCAN 到多表读取的完整配置与源码解析

SeaTunnel Redis Source 连接器实战指南:从单表 SCAN 到多表读取的完整配置与源码解析 SeaTunnel Redis Source 连接器实战指南从单表 SCAN 到多表读取的完整配置与源码解析【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本篇技术指南基于 SeaTunnel 仓库中的 Redis Source 连接器官方文档系统讲解 Redis 数据源连接器的全部配置项、hash 解析模式、键值输出控制与多表读取机制并结合seatunnel-connectors-v2/connector-redis模块的源码如 RedisSourceFactory.java、RedisSourceReader.java逐层剖析其 SCAN 游标扫描、按类型分派读取、多表路由的底层实现帮助你在生产环境中正确完成 Redis 到任意目标端的批量数据抽取。连接器概述与能力边界SeaTunnel 的 Redis Source 连接器插件标识为Redis插件目录为connector-redis见 plugin-mapping.properties 第 73 行seatunnel.source.Redis connector-redis用于从 Redis 中按 key 模式批量读取数据。它的关键能力与限制如下支持引擎Spark、Flink、SeaTunnel Zeta运行模式仅支持batch批处理不支持 stream、exactly-once、列投影与用户自定义 split多表读取支持通过tables_configs一次声明多个 key 模式并读取源码中RedisSource.getProducedCatalogTables()会返回多张CatalogTable见 RedisSource.java。从源码看RedisSource继承自AbstractSingleSplitSource并返回Boundedness.BOUNDED这解释了它为什么只能用于批作业——它把整个 key 扫描任务视为一个有界的单 split。依赖安装使用install-plugin.sh或从 Maven Central 下载connector-redis包放入对应引擎的connector目录即可。核心选项详解连接级选项nametyperequireddefault valueDescriptionhoststringyes when modesingle-Redis server hostportintno6379Redis server portuserstringno-Redis authentication userauthstringno-Redis authentication passworddb_numintno0Redis database indexmodestringnosingleRedis mode:singleorclusternodeslistyes when modecluster-Redis cluster nodes in format[host1:port1, host2:port2]tables_configslistno-List of table configurations for multi-table readingcommon-optionsno-Source plugin common parameters, refer to Source Common Options这些选项在源码中逐一定义于 RedisBaseOptions.javaPORT默认 6379、DB_NUM默认 0、MODE默认SINGLE、BATCH_SIZE默认 10、FORMAT默认JSON、FIELD_DELIMITER默认,与 RedisSourceOptions.java。配置校验规则在 RedisSourceFactory.java 的optionRule()中声明几个值得注意的约束tables_configs与keys互斥OptionRule.builder().exclusive(...)二者只能配置其一当mode single时host必填且非空白port必须在 165535 之间对应源码中的MIN_PORT/MAX_PORT当mode cluster时nodes必填且非空并由RedisNodesValidator校验每个节点是否为合法的host:port格式当read_key_enabled true时single_field_name变为条件必填见conditional(READ_KEY_ENABLED, true, SINGLE_FIELD_NAME)规则。表级选项keys 单表模式与 tables_configs 多表模式通用使用keysdata_type描述单表使用tables_configs描述多表。二者不可同时配置。每张表的配置支持以下参数nametyperequireddefault valuedescriptionkeysstringyes-Redis key pattern to scandata_typestringyes-Redis data type:key,string,hash,list,set,zsetbatch_sizeintno10Batch size for SCAN operationsformatstringnojsonData format:jsonortextschemaconfigno-Schema configuration for this tablehash_key_parse_modestringnoallHash key parse mode:allorkvread_key_enabledbooleannofalseInclude Redis key in outputkey_field_namestringno-Field name for Redis keysingle_field_namestringno-Field name for single-value typesfield_delimiterstringno,Delimiter for text format注意单表模式下这些表级选项直接写在Redis { ... }外层多表模式下它们必须写在tables_configs的每一项内部。多表模式下每项都必须声明非空keys与data_type否则启动即报错——该校验由 RedisTableConfigsValidator.java 执行。从源码结构看单表/多表两种写法最终都汇聚到 RedisTableConfig.java 的of()方法检测到tables_configs时逐项构建否则把顶层配置当作单表构建因此两种写法在运行时共享完全一致的解析逻辑。data_type六种数据类型的读取语义data_type决定每个 key 的值如何拆解成下游行key / string每个 key 的值整体作为一行。例如 key 的值为SeaTunnel test message下游收到一行SeaTunnel test message。hashhash 的所有字段-值对会被格式化为 JSON整体作为一行发送。例如 hash 值为name:tyrantlucifer age:26下游收到一行{name:tyrantlucifer,age:26}。list / set / zset集合中的每个元素单独作为一行发送。例如值为[tyrantlucifer, CalvinKirs]下游收到两行。在 RedisSourceReader.java 的pollNext方法中可以看到这一分派逻辑HASH走pollHashMapToNextSTRING/KEY走pollStringToNextLIST/SET/ZSET分别走各自的 poll 方法其他类型抛出UNSUPPORTED_DATA_TYPE异常。另外一个源码细节data_type key在扫描阶段会被映射为STRING类型执行 SCAN见resolveScanType方法RedisSourceReader.java因为 Redis 的 SCAN TYPE 过滤只有 string/hash/list/set/zstream 等实际类型。重要提示连接器支持模糊 key 匹配glob 通配符用户必须确保同一keys模式匹配到的 key 属于相同的数据类型否则按类型分派读取时会读到错误的值。batch_size 与 SCAN 游标机制batch_size控制每次 SCAN 迭代尝试返回的 key 数量默认 10。底层实现采用 Redis 标准的增量 SCAN 协议从游标0ScanParams.SCAN_POINTER_START开始携带count batch_size与pattern keys调用 SCAN每轮取回一批 key 后立即按data_type分派读取并输出循环直到游标回到0表示整个 keyspace 扫描完毕全部完成后调用context.signalNoMoreElement()通知下游数据源已结束。完整流程见 RedisSourceReader.java 的processTable方法其中还包含进度日志每 100 个 key 或每 10 次迭代打印一次扫描进度结束时输出总 key 数、迭代次数与耗时。SCAN 是非阻塞的因此即使面对大 keyspace 也不会长时间阻塞 Redis 服务端这是它优于 KEYS 命令的关键原因。hash_key_parse_modeall 与 kv 两种解析模式该选项默认all定义于 RedisSourceOptions.java决定 hash 类型 key 的拆解粒度。假设某个 hash key 的值为{ 001: { name: tyrantlucifer, age: 26 }, 002: { name: Zongwen, age: 26 } }模式一hash_key_parse_mode all整个 hash 被当作一行并按 schema 中的字段名逐个子对象解析schema { fields { 001 { name string age int } 002 { name string age int } } }001002Row(nametyrantlucifer, age26)Row(nameZongwen, age26)模式二hash_key_parse_mode kvhash 中的每个 field-value 对被当作一行即行爆炸hash 的 field 名放入 schema 的第一个字段schema { fields { hash_key string name string age int } }hash_keynameage001tyrantlucifer26002Zongwen26源码印证RedisRecordReader.java 的pollHashMapToNext中KV模式会把每个 hash 的字段值映射序列化为 JSON 后交给反序列化 schema 解析而ALL模式或其他模式分支则将整个 hash 序列化为一个 JSON 字符串作为单行的一列输出。注意文档中的提示连接器会用 schema 配置的第一个字段作为每个 hash field 名的输出列名。read_key_enabled把 Redis key 一起带出去默认情况下read_key_enabled false连接器只输出 value。置为true后每条输出记录同时包含 Redis key 与其关联的值典型用途是把缓存中的 key 本身如用户 ID、会话 ID作为业务主键落到目标表。相关配套选项key_field_name[string]key 在输出行中的列名。read_key_enabled true时默认列名为keydata_type hash时若未显式设置默认列名为hash_key对应源码 RedisTableConfig.java 的resolveKeyFieldNamedataType HASH ? hash_key : key。当默认列名与 schema 已有字段冲突时可用它改名例如key_field_name custom_key hash_key_parse_mode kv format json schema { fields { custom_key string name string } }single_field_name[string]当read_key_enabled true且值为单一原始类型如string、int时指定 value 在 schema 中的列名。对可直接映射到 schema 的复杂类型如 hash无效。若配置了 schema务必包含 key 列默认key或key_field_name指定的名称与该列read_key_enabled true key_field_name key single_field_name value schema { fields { key string value string } }从源码看这一开关会决定 Reader 的选择RedisSourceReader.java 的createRecordReader中readKeyEnabled true时创建KeyedRecordReader内部持有KeyValueMerger把 key 合并进待解析的 JSON否则创建UnKeyedRecordReader。且KeyedRecordReader.pollValueToNext强制要求存在反序列化 schema否则抛illegalArgument异常KeyedRecordReader.java——也就是说key-value 联合输出模式下必须配置schema。format 与 schemajson / text 两种上游数据格式format默认json描述从 Redis 取回的值本身的文本形态目前支持json与text。format json必须同时给出schema。例如上游数据是{code: 200, data: get success, success: true}则配置schema { fields { code int data string success boolean } }输出codedatasuccess200get successtrueformat text可以不给 schema——此时上游整串文本会落入默认的单列content行源码中未提供 schema 时构建 simple text table见 RedisTableConfig.java。例如上游数据200#get success#true输出为content200#get success#true若要按分隔符切列则需同时配置schema与field_delimiterfield_delimiter # schema { fields { code int data string success boolean } }底层实现中format决定反序列化 schema 的构建JSON 模式使用JsonDeserializationSchemaTEXT 模式使用带delimiter的TextDeserializationSchema见 RedisTableConfig.javafield_delimiter默认值为,仅在text格式下需要关心。schema 字段语法详见 Schema Feature。完整配置示例示例一最简单的 key 读取Redis { host localhost port 6379 keys key_test* data_type key format text }示例二按 schema 解析 JSON 值Redis { host localhost port 6379 keys key_test* data_type key format json schema { fields { name string age int } } }示例三读取 string 类型 key 并追加写入 Redis list该示例同时展示了 Source 与 Sink 两侧的 Redis 配置Source 侧使用keys模糊匹配 batch_sizeSink 侧把结果写入固定 key 的 listsource { Redis { host redis-e2e port 6379 auth U2VhVHVubmVs keys string_test* data_type string batch_size 33 } } sink { Redis { host redis-e2e port 6379 auth U2VhVHVubmVs key string_test_list data_type list batch_size 33 } }示例四带 key 联合输出的 string 读取source { Redis { host redis-e2e port 6379 auth U2VhVHVubmVs keys string_test* data_type string batch_size 33 read_key_enabled true key_field_name custom_key single_field_name custom_value format json schema { table RedisDatabase.RedisTable columns [ { name custom_key type string }, { name custom_value type string } ] } } } sink { Console {} }多表模式tables_configs示例五一次读取多种 key 模式与数据类型env { job.mode BATCH } source { Redis { host localhost port 6379 auth password tables_configs [ { keys user:active:* data_type STRING format JSON batch_size 10 schema { table user_table fields { id int name string email string created_at timestamp } } }, { keys session:* data_type HASH hash_key_parse_mode KV read_key_enabled true key_field_name session_id schema { table session_table fields { session_id string user_id int ip_address string last_active timestamp } } }, { keys queue:task:* data_type LIST format TEXT field_delimiter | schema { table task_table fields { content string } } } ] } } sink { Redis { host localhost port 6379 key redis-result-${table_name} data_type LIST } }多表模式下的表名来源于各表schema.table源码中由getTablePath方法从 schema 配置读取TableIdentifierOptions.TABLE缺省时回退为 key 模式本身见 RedisTableConfig.java。下游 sink 可使用${table_name}占位符把不同 key 模式的数据路由到不同的目标表或 key。源码层面有两个值得注意的机制每条输出行都会通过setTableId(tablePath.toString())打上表标识见 RedisRecordReader.java这是${table_name}路由能够生效的基础RedisSource.createSourceTablesMap会对重复的表路径直接抛错Duplicate table_path foundRedisSource.java因此多表模式下各表的schema.table必须唯一。示例六集群模式 多表source { Redis { mode CLUSTER nodes [node1:6379, node2:6379, node3:6379] auth cluster_password tables_configs [ { keys metric:cpu:* data_type STRING format JSON batch_size 10 schema { fields { host string timestamp timestamp usage double } } }, { keys metric:memory:* data_type STRING format JSON batch_size 10 schema { fields { host string timestamp timestamp used long total long } } } ] } } sink { Console {} }源码视角连接构建与版本探测集群/单实例的连接构建集中在 RedisParameters.javasingle 模式创建 Jedis 单连接依次执行auth配置了auth时、用户认证配置了user时与select(db_num)切换到指定库再包装为RedisSingleClientcluster 模式把nodes解析为host:port集合创建JedisCluster配置了auth时传入密码包装为RedisClusterClient版本探测buildRedisClient会先调用INFO命令解析redis_version如5.0.14取主版本号 5解析失败抛出GET_REDIS_VERSION_INFO_FAILED错误RedisParameters.java——从源码结构看这个版本号会被存入RedisClient用于在客户端层面适配不同 Redis 版本的行为差异。测试覆盖方面仓库提供了针对 Redis 5 与 Redis 7 的测试Redis5Test.java、Redis7Test.java以及多表配置解析的单元测试 RedisTableConfigTest.javaE2E 层则有一组 Redis 作业配置可参考例如 redis-to-redis.conf 与集群节点格式非法时的校验用例 redis-validation-cluster-malformed-node.conf。常见配置问题排查清单基于上文文档约束与源码校验逻辑配置报错时可对照以下清单tables_configs与keys同时出现OptionRule 声明二者互斥二者只保留其一多表模式下某项缺keys或data_type会抛出tables_configs[i]: keys must be configured and non-blank一类错误mode cluster但nodes格式错误必须是[host1:port1, ...]每个元素按第一个:拆分为 host 与 portmode single但端口超出 1~65535会被条件校验拦截E2E 用例 redis-validation-single-port-out-of-range.conf 即验证此场景read_key_enabled true却未配single_field_name/schema单值类型下single_field_name为条件必填且 schema 必须包含 key 列与 value 列模糊 key 匹配到了不同类型同一keys模式下的 key 必须同类型否则读取语义错乱。小结SeaTunnel Redis Source 连接器是一个纯批处理的 Redis 全量抽取组件以 SCAN 增量游标 batch_size控制扫描压力用data_type区分六种 Redis 数据结构的行化语义用hash_key_parse_mode、read_key_enabled/key_field_name/single_field_name精细控制 key 与 value 的输出形态并通过tables_configsschema.table${table_name}实现一次作业多模式、多类型的抽取与路由。掌握上述选项语义后再结合connector-redis模块的源码工厂校验、参数构建、Reader 分派三层可以准确判断任意一份 Redis 抽取配置的实际行为。参考文档Redis Source 官方文档、Source Common Options、Schema Feature。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表