ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

SeaTunnel Typesense Sink 连接器实战指南:写入模式、文档 ID 生成与批量刷新机制

2026/9/18 22:49:40 拓冰建站 浏览量
SeaTunnel Typesense Sink 连接器实战指南:写入模式、文档 ID 生成与批量刷新机制 SeaTunnel Typesense Sink 连接器实战指南写入模式、文档 ID 生成与批量刷新机制【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本篇指南围绕 SeaTunnel 的 Typesense Sink 连接器 展开系统讲解如何将 SeaTunnel 数据写入 Typesense collection从连接参数、schema_save_mode/data_save_mode两种写入前处理策略到基于primary_keys的文档id生成、批量写入与重试机制并给出批处理、读写对拷与流式 upsert 三套可直接运行的 HOCON 作业配置。读完你可以独立完成 Typesense 相关的 SeaTunnel 数据集成任务并理解其底层写入路径。连接器概述Typesense Sink 是 SeaTunnel Connector V2 体系中的写入端插件连接器标识为Typesense定义于 TypesenseBaseOptions.java。它负责把 SeaTunnel 内部行数据SeaTunnelRow序列化为 Typesense 文档并批量写入目标 collection。连接器具备以下能力按配置自动创建目标 collection支持基于上游表结构建表写入前清理已有文档DROP_DATA使用一个或多个主键字段生成稳定的 Typesense 文档id实现 upsert 语义支持多表写入路由到同一 collection。支持的引擎SeaTunnel Zeta主要特性特性支持情况精确一次exactly-once❌ 不支持CDC✅ 支持多表写入✅ 支持批处理✅ 支持流处理✅ 支持定时刷新❌ 不支持关于这些特性的通用定义可参考 Connector V2 功能简介。连接器在流式模式下依赖 checkpoint 触发批量落盘因此不支持独立于 checkpoint 的定时刷新。连接参数与选项详解以下是 Typesense Sink 选项 中定义的完整参数表默认值与类型均与源码 TypesenseSinkOptions.java 及 TypesenseBaseOptions.java 保持一致名称类型是否必须默认值描述hostsarray是-Typesense 节点地址格式为host:port支持配置多个地址collectionstring是-目标 collection 名schema_save_modestring是CREATE_SCHEMA_WHEN_NOT_EXIST写入前如何处理目标 collection 结构data_save_modestring是APPEND_DATA写入前如何处理目标 collection 中已有文档primary_keysarray否-用于生成 Typesense 文档id的源字段key_delimiterstring否_primary_keys配置多个字段时使用的拼接分隔符api_keystring是-Typesense API Keymax_retry_countint否3单个批量请求的最大重试次数max_batch_sizeint否10单个批量请求最多写入的文档数量multi_table_sink_replicaint否1通用多表写入路由机制使用的 Sink 副本数common-options否-通用 Sink 选项hosts [array]Typesense 的访问地址格式为host:port例如[typesense-01:8108]。从源码 TypesenseClient.createInstance 可以看出配置的每个 host 都会被解析为Node(protocol, host, port)加入节点列表若host:port中未显式给出端口端口部分为空会回退到默认端口8018。配置多个节点时每个 Writer 只持有一个客户端不会把写入请求在节点之间做负载均衡。collection [string]要写入的 collection 名例如seatunnel。在多表作业中所有表都会路由到同一个 collection这一点在 FixedValueCollectionSerializer.java 的实现中可以确认如果不同表要写入不同目标请为每个目标 collection 单独配置一个 sink 块。primary_keys [array]主键字段用于生成文档id。配置多个字段时连接器会用key_delimiter拼接这些字段值。未配置primary_keys时Typesense 会自行分配文档 ID连接器退化为纯追加写入。底层实现见 KeyExtractor.java当primary_keys为null时createKeyExtractor返回row - null随后 TypesenseRowSerializer.serializeRow 在生成文档时不会写入id字段交由 Typesense 服务端自动分配。key_delimiter [string]设定复合键的分隔符默认为_。源码注释中的示例设置为$时文档id形如KEY1$KEY2$KEY3。api_key [string]Typesense 安全认证的api_key。请把它当作敏感凭据处理在共享环境运行时建议通过作业密钥或环境变量注入。max_retry_count [int]单个批量请求的最大重试次数。在 TypesenseSinkWriter 中重试谓词为exception - true也就是说typesenseClient.insert(...)抛出的任何异常网络错误、超时以及 Typesense 业务错误响应都会被同样重试最多执行max_retry_count次每次间隔固定的 200 msDEFAULT_SLEEP_TIME_MS 200L当前实现并不会区分瞬时错误和永久错误。重试仍失败时会抛出INSERT_DOC_ERROR对应的连接器异常终止任务。max_batch_size [int]每批最多写入的文档数量。Writer 会在批次达到max_batch_size时立即触发一次批量请求在 checkpoint 或关闭时再冲刷剩余数据。Typesense 对单次请求有上限请将该值保持在 Typesense 服务端per_page上限以下。multi_table_sink_replica [int]通用多表写入选项。当多表任务需要为 Typesense 写入端配置更多 Sink 副本时使用。common optionsSink 插件常用参数plugin_input、parallelism、metadata_datasource_id等请参考 Sink 常用选项 了解详情。其中plugin_input用于指定上游数据集当不指定时当前插件处理配置文件中上一个插件输出的数据集dataset当指定时处理该参数对应的数据集。schema_save_mode 与 data_save_mode写入前的双阶段处理schema_save_mode与data_save_mode是 SeaTunnel Sink 侧通用的“写入前处理”机制。Typesense Sink 实现了SupportSaveMode接口通过 TypesenseSink.getSaveModeHandler 将两个模式与TypesenseCatalog一起封装进DefaultSaveModeHandler在任务启动同步之前执行对应的建表、删表、清数据等动作。schema_save_mode在启动同步任务之前针对目标侧已有的表结构选择不同的处理方案。选项介绍RECREATE_SCHEMA当表不存在时会创建当表已存在时会删除并重建CREATE_SCHEMA_WHEN_NOT_EXIST当表不存在时会创建当表已存在时则跳过创建默认值ERROR_WHEN_SCHEMA_NOT_EXIST当表不存在时将抛出错误Typesense collection 创建时会使用上游 SeaTunnel 表结构。如果希望重复写入时文档id保持稳定请配置primary_keys。从 TypesenseClient.createCollection 与 TypesenseCatalog.createTable 可以看到建表动作实际调用 Typesense 原生 SDK 创建 collection默认字段为.*类型AUTO并开启enableNestedFields(true)即允许嵌套对象字段删除重建则对应dropCollection。truncateTable对应 truncateCollectionData内部通过filterBy(id:!1||id:1)的删除参数清空全部文档。data_save_mode在启动同步任务之前针对目标侧已存在的数据选择不同的处理方案。选项介绍DROP_DATA保留数据库结构删除数据APPEND_DATA保留数据库结构保留数据默认值ERROR_WHEN_DATA_EXISTS当有数据时抛出错误:::tip连接器使用 Typesense 的批量导入接口且导入参数中固定设置了action(upsert)见 TypesenseClient.insert。UPDATE和DELETE行类型不会被解释为 CDC 操作 —— 每条上游记录都会按生成的文档id被 upsert 到目标 collection。如果希望重复作业行为类似 upsert 而不是追加可以把data_save_mode设为DROP_DATA并配置稳定的primary_keys。:::文档 ID 生成与行序列化原理写入路径的核心是把SeaTunnelRow变成一条 Typesense 文档 JSON。整个链路为TypesenseRowSerializer.serializeRow 先通过KeyExtractor计算主键再调用toDocumentMap把行按字段名转成 Map若 key 非空则写入id字段最后由 JacksonObjectMapper序列化为 JSON 字符串。toDocumentMap会递归处理嵌套结构字段值如果是SeaTunnelRowROW 类型会递归转成嵌套 Map因此上游表结构中的嵌套对象可以原样保留为 Typesense 的嵌套字段配合建表时的enableNestedFields(true)生效。convertValue对TemporalJDK 8 时间类型如LocalDateTime、LocalDate、LocalTime统一调用toString()转为字符串注释说明 jackson 不支持 JDK8 新时间 API并对Map、List内部元素递归转换。主键提取方面KeyExtractor 对每个主键字段按顺序取值并用key_delimiter拼接成最终id字符串。需要注意ROW、ARRAY、MAP 类型不能作为主键字段会抛出UNSUPPORTED_OPERATION异常DATE / TIME / TIMESTAMP 类型的主键会转成字符串形式参与拼接。此外写入端对行的RowKind有专门的分流处理TypesenseSinkWriter.writeINSERT/UPDATE_AFTER序列化为文档 JSON追加进待写入批次UPDATE_BEFORE直接忽略不产生任何写入DELETE通过serializeRowForDelete提取id调用deleteCollectionData(collection, id)按文档 ID 删除单条记录其他行类型抛出UNSUPPORTED_OPERATION异常。批量写入、重试与 Checkpoint 刷新机制Writer 的批量行为是理解该连接器吞吐与一致性模型的关键逻辑全部集中在 TypesenseSinkWriter.java触发时机一按量刷新write()每累积一条INSERT/UPDATE_AFTER记录就检查requestEsList.size() maxBatchSize达到阈值立即调用insert(collection, requestEsList)发出批量请求成功后清空缓冲列表。触发时机二checkpoint 刷新prepareCommit()在每次 checkpoint 前被调用把缓冲区中剩余的全部记录一次性插入。触发时机三关闭刷新close()同样会冲刷剩余数据保证任务结束时没有残留。重试每次批量插入都包在RetryUtils.retryWithException中重试次数由max_retry_count控制默认 3每次重试固定间隔 200 ms重试谓词对所有异常恒为true不区分瞬时/永久错误。这意味着在流式模式下Writer 最多缓冲max_batch_size条记录或者直到下一个 checkpoint才发出一次批量请求。把data_save_mode DROP_DATA与稳定的primary_keys组合起来每个 checkpoint 都会产生幂等的 upsert —— 这正是本文第三个任务示例所展示的“流式 Upsert 并按 Checkpoint 刷新”模式。任务示例以下示例均取自 Typesense Sink 文档并保留完整的可运行配置。使用主键写入文档批处理模式下用FakeSource生成 5 行数据primary_keys指定num_employees与num两个字段用拼接成文档idenv { parallelism 1 job.mode BATCH } source { FakeSource { row.num 5 plugin_output typesense_test_table schema { fields { company_name string num long id string num_employees int flag boolean } } } } sink { Typesense { plugin_input typesense_test_table hosts [localhost:8108] collection typesense_test_collection api_key xyz primary_keys [num_employees, num] key_delimiter max_retry_count 3 max_batch_size 10 schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }从 Typesense 读取并写入另一个 collection利用 Typesense Source 与 Sink 搭配在一个作业内完成 collection 到 collection 的复制。Source 端通过query传入 Typesense 搜索表达式q*filter_byc_row.c_int:10并显式声明schema包含嵌套对象c_rowSink 端用num_employees与id拼接主键env { parallelism 1 job.mode BATCH } source { Typesense { hosts [localhost:8108] collection typesense_source_collection api_key xyz query q*filter_byc_row.c_int:10 plugin_output typesense_test_table schema { fields { company_name_list arraystring company_name string num_employees long country string id string c_row { c_int int c_string string c_array_int arrayint } } } } } sink { Typesense { plugin_input typesense_test_table hosts [localhost:8108] collection typesense_sink_collection api_key xyz primary_keys [num_employees, id] key_delimiter max_retry_count 3 max_batch_size 10 schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }流式 Upsert 并按 Checkpoint 刷新流式模式下parallelism 2、checkpoint.interval 3000030 秒一个 checkpoint。Writer 最多缓冲max_batch_size条记录或者直到下一个 checkpoint再发出一次批量请求。把data_save_mode DROP_DATA与稳定的primary_keys组合起来每个 checkpoint 都会产生幂等的 upsertenv { parallelism 2 job.mode STREAMING checkpoint.interval 30000 } source { FakeSource { row.num 1000 schema { fields { company_name string num long id string num_employees int flag boolean } } plugin_output typesense_stream } } sink { Typesense { plugin_input typesense_stream hosts [localhost:8108] collection typesense_stream_collection api_key xyz primary_keys [id] max_batch_size 100 schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode DROP_DATA } }深入源码连接器模块结构若想进一步研究实现细节可在仓库seatunnel-connectors-v2/connector-typesense模块下按以下路径阅读配置定义TypesenseSinkOptions.java、TypesenseBaseOptions.java含未写入文档但默认生效的protocol选项默认http官方注释建议 Typesense Cloud 场景使用https客户端封装TypesenseClient.java建表、删表、清空、按 id 删除、文档数统计、批量 upsert 导入写入实现TypesenseSinkWriter.java、TypesenseSink.java序列化TypesenseRowSerializer.java、KeyExtractor.javaCatalog / Save Mode 落地TypesenseCatalog.java测试用例TypesenseRowSerializerTest.java、TypesenseFactoryTest.java可作为理解序列化与工厂注册行为的参考变更日志该连接器的历史变更记录见 connector-typesense 变更日志。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考