ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

SeaTunnel Typesense Source 连接器:从 Typesense 集合批量抽取数据的配置与源码解析

SeaTunnel Typesense Source 连接器:从 Typesense 集合批量抽取数据的配置与源码解析 SeaTunnel Typesense Source 连接器从 Typesense 集合批量抽取数据的配置与源码解析【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文基于 SeaTunnel 仓库中 Typesense Source 连接器的官方文档Typesense.md与对应源码实现完整讲解该连接器在 Zeta 引擎下批量读取 Typesense 集合文档的配置参数、作业示例与分页拉取原理。读完本文你可以独立完成一个 Typesense 到 SeaTunnel 下游Console、JDBC 等的批量同步作业并理解其单 Split 枚举、query参数透传与offset/per_page分页的底层机制。一、连接器定位有界批量读取不做 CDCTypesense Source 连接器connector 标识为Typesense见 TypesenseBaseOptions.java 中的CONNECTOR_IDENTITY用于从 Typesense 的某个 collection 中读取文档。它的核心定位是支持引擎仅 SeaTunnel Zeta 引擎有界bounded数据源每个作业只会把符合query条件的文档完整读取一遍读完后即结束。从源码看TypesenseSource.java 的getBoundedness()直接返回Boundedness.BOUNDED不支持流式与 CDC如果需要同步 Typesense 的增量变更官方建议由接收端使用 Typesense 自带的 change tracking 机制配合消费而不是依赖本连接器支持的能力Batch Processing、Schema 声明、Parallelism支持并行度配置但见下文原理说明其扫描实际为单 Split 模型不支持 Exactly-Once 与用户自定义 Split。适用场景举例将存量文档从 Typesense 全量/按条件导出到数仓或对象存储定期用filter_by过滤出特定子集做批量加工。二、Options 参数详解完整参数表继承自官方文档NameTypeRequiredDefaulthostsarrayyes-collectionstringyes-schemaconfigyes-api_keystringyes-protocolstringnohttpquerystringno-batch_sizeintno100common-optionsno-这些参数在源码中的定义分为两处基础参数位于 TypesenseBaseOptions.javasource 专属参数位于 TypesenseSourceOptions.java二者的默认值与必填性与文档描述一致protocol默认http、batch_size默认100QUERY无默认值。hosts [array]Typesense 的访问地址格式为host:port例如[typesense-01:8108]支持配置多个节点。需要特别注意两点源码行为配置多个节点时source 的搜索请求会发送到第一个可达节点该列表不会用于并行扫描分片——从源码结构看TypesenseSourceSplitEnumerator.java 的getTypesenseSplit()约 L145-L156每个作业只生成一个TypesenseSourceSplit因此扫描本身是单线程拉取的如果host:port中省略端口TypesenseClient.java 的createInstance()L78-L97会将端口兜底为8018并统一使用protocol参数构建 Node。collection [string]要读取的 Typesense collection 名称例如companies。schema [config]声明要从 Typesense 读取的列及其类型写法参见 Schema 特性指南其中 How to declare type supported 一节。source 侧通过 TypesenseSource.java 中的CatalogTableUtil.buildWithConfig(config)将schema构建成CatalogTable再传给 Reader 作为行反序列化的SeaTunnelRowType并实现SupportColumnProjection支持列投影。api_key [string]Typesense 的 API Key用于安全认证。该值属于敏感凭证在共享基础设施上运行时建议通过作业密钥或环境变量传入而不是硬编码在配置文件中。protocol [string]连接 Typesense 使用的协议默认http。对接 Typesense Cloud 或其他启用 TLS 的端点时使用https。query [string]Typesense 搜索参数例如q*filter_bynum_employees:9000。不配置时source 使用 Typesense 默认搜索返回的全部文档。任何合法的 Typesense 搜索参数都可以追加包括q、query_by、filter_by、sort_by、page、per_page连接器将它们原样透传给 Typesense 搜索 API参数格式为 URL 查询串风格k1v1k2v2。从源码看URLParamsConverter.java 的parseParams()L42-L61按分割、再按最多切两段解析成键值对任何一个片段缺少都会抛出QUERY_PARAM_ERROR异常随后在 TypesenseClient.java 的search()L121-L136中反序列化为SearchParameters对象。若query为空则退化为默认q(*)全量搜索。batch_size [int]每次批量读取的记录数默认100。每次请求实际使用 Typesense 的per_page参数因此取值必须介于 1 与 Typesense 服务端per_page上限通常为 250之间。如果在日志中看到分页被截断应调小该值。Common OptionsSource 插件的通用参数如source.schema之外的公共项参见 Source Common Options。三、端到端工作原理从 Split 枚举到分页拉取结合源码可以还原 Typesense Source 的完整执行链路作业启动与 Split 枚举Zeta 引擎调用TypesenseSource.createEnumerator()创建 TypesenseSourceSplitEnumerator.java。其run()方法中首次枚举时读取collection、query、batch_size三个配置构造一个携带SourceCollectionInfocollection 名、query、初始 offset0、batch size的 Split 并分发给 Reader随后向所有 Reader 发送NoMoreSplitsEvent。枚举器还支持通过snapshotState()保存shouldEnumerate、pendingSplit与assignCount状态配合 TypesenseSourceState.java 实现断点续跑时的状态恢复。Reader 创建客户端Reader 打开时通过TypesenseClient.createInstance(config)用hosts/protocol/api_key构建 Typesense Java 客户端连接超时 5 秒见 TypesenseClient.java L94-L95。分页搜索循环TypesenseSourceReader.java 的pollNext()L92-L128执行核心拉取逻辑每次调用typesenseClient.search(collection, query, offset, pageSize)即带offset与per_page的搜索请求遍历返回的hits将每个文档 Map 经DefaultSeaTunnelRowDeserializer按schema反序列化为SeaTunnelRow后输出用响应中的found判断是否还有下一页当found / pageSize - 1 offset / pageSize时offset增加pageSize继续翻页否则结束当前 SplitSplit 全部消费完且收到noMoreSplit信号后Reader 上报signalNoMoreElement()作业正常收尾。这一结构解释了前文两个实践要点多 host 列表只影响连接容错而不影响并行度batch_size直接决定每页大小并参与终止条件判断超范围会导致服务端截断页面。四、作业配置示例示例一按过滤条件读取文档query中使用filter_by只抽取员工数大于 9000 的公司env { parallelism 1 job.mode BATCH } source { Typesense { hosts [localhost:8108] collection companies api_key xyz query q*filter_bynum_employees:9000 batch_size 100 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 { Console {} }注意job.mode BATCH是必须的因为该 source 是 BOUNDED 的schema 中c_row展示了嵌套对象字段的声明方式。示例二自定义 query_by 与 sort_by 读取子集组合query_by与sort_by来控制 Typesense 在哪些字段上搜索、结果集如何排序env { parallelism 1 job.mode BATCH } source { Typesense { hosts [localhost:8108] collection companies api_key xyz query qacmequery_bycompany_namefilter_bycountry:USsort_bynum_employees:desc batch_size 50 schema { fields { company_name string num_employees long country string id string } } } } sink { Console {} }该示例中query串包含四个参数在company_name字段上全文搜索 acme过滤美国的公司并按员工数降序输出。五、使用限制与注意事项必须 BATCH 模式连接器仅有界不支持流式消费需要增量同步时应在下游消费 Typesense 自身的 change tracking 事件batch_size上限受 Typesense 服务端per_page限制典型值 250超限或日志出现页面截断时应调小单 Split 扫描模型从 TypesenseSourceSplitEnumerator.java 的getTypesenseSplit()结构看作业只生成一个 Split提升parallelism不会把 Typesense 扫描本身打散成多路并行query 语法要求query内每个参数必须是keyvalue形式片段数量以分隔否则初始化搜索参数时即抛错对应测试见 URLParamsConverterTest.java凭证安全api_key应按敏感信息管理共享环境中优先走环境变量或密钥通道。六、版本演进根据仓库内的变更日志connector-typesense.mdTypesense 连接器于 2.3.8 版本首次引入#7450随后在 2.3.9 至 2.3.12 版本中持续完善包括 shade 检查规则、枚举器 API 语义优化#9671与 options 结构改进#9398等。使用新版 SeaTunnel2.3.8 及以上即可获得文档中描述的全部参数能力。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表