消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Elasticsearch sink connector 是 Apache Pulsar 官方提供的 IO 连接器之一它从 Pulsar topic 中拉取消息并持久化到 Elasticsearch 索引中。本文以仓库文档 site2/website-next/docs/io-elasticsearch-sink.md 为主体结合连接器源码位于 pulsar-io/elastic-search深入讲解其 Raw / Schema 感知两种数据处理模式、_id主键映射、多索引写入、Bulk 批量写入与 TLS 安全连接并给出完整可运行的配置与启动示例。读完本文你将掌握该连接器的全部配置项含义、参数默认值与校验规则以及从零部署一套 Elasticsearch sink 的完整操作路径。连接器概述与整体工作流程Elasticsearch sink connector 以SinkGenericObject的形式实现通过write()方法将每条 Pulsar 记录转换为一个 Elasticsearch 文档请求index / delete再交给底层 REST 客户端执行。从源码结构看连接器核心由以下类组成类相对路径pulsar-io/elastic-search/src/main/java/org/apache/pulsar/io/elasticsearch/职责ElasticSearchSink.java连接器主体负责打开连接、提取_id与_source、分发写入请求ElasticSearchClient.java底层 REST 客户端封装含 Bulk 处理器、重试策略、TLS、索引自动创建ElasticSearchConfig.java全部配置项定义与validate()校验逻辑IndexNameFormatter.java基于事件时间event time的索引名格式化JsonConverter.java将 AVRO GenericRecord 及其 logical type 转换为 JSON 节点连接器通过Connector(name elastic_search, type IOType.SINK, configClass ElasticSearchConfig.class)注解注册见 ElasticSearchSink.java构建后产出 NAR 归档文件pulsar-io-elastic-search-version.nar由 Maven 的nifi-nar-maven-plugin打包见 pulsar-io/elastic-search/pom.xml。两种数据处理模式Raw processing 与 Schema aware自 Pulsar 2.9.0 起Elasticsearch sink connector 提供两种工作方式通过schemaEnable参数切换模式说明Raw processingSink 从 topic 读取消息并把原始内容直接写入 Elasticsearch。这是默认行为schemaEnable默认为false。该模式在 Pulsar 2.8.x 中已经可用。Schema awareSink 使用消息的 Schema 解析内容支持 AVRO、JSON 和 KeyValue 三种 Schema 类型并把内容映射为 Elasticsearch 文档。当schemaEnable设为true时连接器会解释消息内容你可以通过primaryFields定义主键该主键会被用作 Elasticsearch 文档的特殊_id字段。Schema aware 模式下的增删改操作Schema aware 模式允许你基于消息的逻辑主键对 Elasticsearch 执行UPDATE、INSERT和DELETE操作。这在典型的Change Data CaptureCDC场景中非常有用数据库变更通过 Debezium 适配器写入 Pulsar再经此 sink 写入 Elasticsearch实现数据库变更的实时同步。通过primaryFields配置项定义主键到_id的映射。DELETE操作当主键不为空且其余值为空时触发删除。用nullValueAction配置该行为默认配置是直接忽略此类空值IGNORE。从源码看ElasticSearchSink.extractIdAndDocument() 是核心逻辑若schemaEnabletrue先判断 record 的 schema 是否为KeyValueSchema若是则拆出 key/value 两套 schema 与值否则以 record key 作为 key、record value 作为 value。若keyIgnorefalse且 key 与 keySchema 均存在则通过stringifyKey()将 key 序列化为_id支持 INT8/INT16/INT32/INT64/STRING/JSON/AVRO 类型。若配置了primaryFields则会从 value 序列化出的 JSON 文档中提取对应字段拼装_id单字段时直接转换为字符串多字段时生成字段值 JSON 数组的字符串表示见 stringifyKey()。若schemaEnablefalseRaw 模式则直接把消息字节按 UTF-8 解码作为文档_source_id为null交由 Elasticsearch 自动生成随机文档 ID。在 write() 中当文档值为空null时按nullValueAction分发DELETE则执行deleteDocument/bulkDeleteIGNORE则跳过FAIL则标记为不可恢复错误并抛异常。当 JSON 序列化失败时按malformedDocActionIGNORE/WARN/FAIL处理。AVRO logical type 的 JSON 转换在 Schema aware 模式下AVRO 记录由 JsonConverter.java 转换为 JSON。它内置了decimal、date、time-millis、time-micros、timestamp-millis、timestamp-micros、uuid等 logical type 转换器见 JsonConverter.java确保 AVRO 中的特殊类型能正确落到 Elasticsearch 文档中。多索引映射indexName 不再必填自 Pulsar 2.9.0 起indexName属性不再是必填项。如果省略sink 会使用 Pulsar topic 名派生索引名。从 ElasticSearchClient.indexName() 的源码可以看到完整逻辑若配置了indexName则使用配置值配合IndexNameFormatter支持时间格式。否则取 record 的 topic 名经 topicToIndexName() 转换先转小写再剥掉persistent://tenant/namespace/前缀只保留最后一段随后截断到 255 字节以内并校验其符合 Elasticsearch 索引命名规则。另外indexName支持日期格式以生成基于事件时间的索引模式为%{date-format}。例如记录的 event time 为1645182000000L即 2022-02-18indexName设为logs-%{yyyy-MM-dd}则实际索引名为logs-2022-02-18。实现位于 IndexNameFormatter.java它用正则%\{\(.?)}解析格式串并用DateTimeFormatter配合 UTC 时区格式化事件时间若配置了createIndexIfNeeded还会校验索引名只能使用小写字符见 validate()。Bulk 批量写入自 Pulsar 2.9.0 起可通过bulkEnabledtrue启用批量写入。底层使用 Elasticsearch Java REST 客户端的BulkProcessor见 ElasticSearchClient.java可以按请求数量、请求大小或时间间隔三种维度触发 flushbulkActions每个 bulk 请求的最大 action 数默认 1000-1 禁用。bulkSizeInMbbulk 请求的最大体积默认 5MB-1 禁用。bulkFlushIntervalInMs最大等待 flush 时间默认 -1 即不设置。bulkConcurrentRequests在途 bulk 请求的最大并发数。默认 0 表示只允许执行单个请求设为 1 表示在累积新 bulk 请求的同时允许 1 个并发请求执行。需要注意如果 pending 的 bulk 请求超过bulkConcurrentRequests下一个 bulk 请求会阻塞即connector.write()会阻塞内存中最多保留(bulkConcurrentRequests 1) * bulkActions条记录该说明见 ElasticSearchConfig.java。Bulk 模式下每条记录会在afterBulk回调中根据BulkItemResponse的结果执行record.ack()或record.fail()见 ElasticSearchClient.java从而实现 Pulsar 的 at-least-once 语义。flush()与close()时都会等待 bulk 处理器清空awaitClose(5000L, MILLISECONDS)。通过 TLS 建立安全连接自 Pulsar 2.9.0 起可以启用 TLS 加密通信。通过ssl子配置块ElasticSearchSslConfig实现。从 ElasticSearchClient.ConfigCallback 的源码可以看到当ssl.enabledtrue时连接管理器会注册https协议加载 truststore/keystore 构建SSLContext并可按需关闭主机名校验hostnameVerificationfalse时使用NoopHostnameVerifier。配置属性详解连接器配置以ElasticSearchConfig类承载ElasticSearchConfig.java通过 Jackson 从 JSON/YAML 加载并在open()时执行validate()校验。以下为完整属性表名称类型必填默认值说明elasticSearchUrlString是连接器连接的 Elasticsearch 集群 URL。可配置多个地址以英文逗号分隔源码中按,拆分并逐一解析为HttpHost。indexNameString否2.9.0 起写入消息的索引名默认值为 topic 名。支持%{date-format}日期格式以生成基于事件时间的索引。typeNameString否_doc写入消息的 type 名。Elasticsearch 6.2 之前的版本需显式设置为除_doc之外的有效 type 名其他版本保持默认。schemaEnableBoolean否false开启 Schema aware 模式。createIndexIfNeededBoolean否false索引缺失时自动创建。indexNumberOfShardsint否1索引分片数校验要求 0。indexNumberOfReplicasint否0文档表中为 1源码默认值为 0索引副本数校验要求 0。注意文档属性表写的是默认 1但当前仓库源码 ElasticSearchConfig.java 中实际默认值为0以源码为准。maxRetriesInteger否1Elasticsearch 请求的最大重试次数设为 -1 禁用。retryBackoffInMsInteger否100重试请求时的基础等待时间毫秒。maxRetryTimeInSecInteger否86400重试请求的最大时间间隔秒。bulkEnabledBoolean否false启用 Elasticsearch bulk processor 批量写入。bulkActionsInteger否1000每个 bulk 请求的最大 action 数-1 禁用。bulkSizeInMbInteger否5bulk 请求的最大体积MB-1 禁用。bulkConcurrentRequestsInteger否0在途 bulk 请求的最大并发数。bulkFlushIntervalInMsInteger否-1启用 bulk 后等待 flush 的最大时间毫秒-1 表示不设置。compressionEnabledBoolean否false启用 Elasticsearch 请求压缩。connectTimeoutInMsInteger否5000客户端连接超时毫秒。connectionRequestTimeoutInMsInteger否1000从连接池获取连接的超时毫秒。connectionIdleTimeoutInMsInteger否30000文档表中为 5空闲连接超时防止连接因空闲被服务端断开而触发读超时。注意文档属性表写的是 5当前仓库源码 ElasticSearchConfig.java 中默认值为30000毫秒以源码为准。socketTimeoutInMsInteger否60000等待读取响应的 socket 超时毫秒。keyIgnoreBoolean否true是否忽略记录 key 来构建文档_id。若定义了primaryFields则从 payload 中提取主键字段构建_id若未提供则由 Elasticsearch 自动生成随机_id。primaryFieldsString否用逗号分隔、有序排列的字段名列表用于从记录 value 构建文档_id。若列表只有一个字段该字段值转为字符串若有两个及以上字段生成的_id为字段值 JSON 数组的字符串表示。nullValueActionenumIGNORE/DELETE/FAIL否IGNORE处理值为 null 的记录。默认 IGNORE 忽略该消息。malformedDocActionenumIGNORE/WARN/FAIL否FAIL处理因格式错误被 Elasticsearch 拒绝的文档。默认 FAIL。stripNullsBoolean否true为 false 时_source对空字段保留null如{foo: null}为 true 时剥离空字段。usernameString否连接 Elasticsearch 集群的用户名。设置后必须同时提供password。passwordString否连接密码。设置后必须同时提供username。sslElasticSearchSslConfig否-TLS 加密通信配置。validate() 中的校验规则ElasticSearchConfig.validate()见 ElasticSearchConfig.java会在连接器启动时校验elasticSearchUrl不能为空否则抛出IllegalArgumentException(elasticSearchUrl not set.)。当indexName与createIndexIfNeeded同时配置时索引名必须全部小写不能以-/_/开头不能是.或..UTF-8 字节长度不能超过 255。username与password必须成对出现二者其一为空则报错。indexNumberOfShards必须为正整数indexNumberOfReplicas必须非负。connectTimeoutInMs、connectionRequestTimeoutInMs、socketTimeoutInMs、bulkConcurrentRequests必须非负。ElasticSearchSslConfig 结构定义名称类型必填默认值说明enabledBoolean否false启用 SSL/TLS。providerString否-SSL 提供方源码中支持如 JCE 提供者名。hostnameVerificationBoolean否true使用 SSL 时是否校验节点主机名。truststorePathString否truststore 文件路径。truststorePasswordString否truststore 密码。keystorePathString否keystore 文件路径。keystorePasswordString否keystore 密码。cipherSuitesString否SSL/TLS 密码套件逗号分隔。protocolsString否TLSv1.2启用的 SSL/TLS 协议列表逗号分隔。配置示例使用连接器前需先通过 JSON 或 YAML 方式创建配置文件。仓库测试目录下也提供了一个参考配置 pulsar-io/elastic-search/src/test/resources/sinkConfig.yaml。Elasticsearch 6.2 之后JSON{ configs: { elasticSearchUrl: http://localhost:9200, indexName: my_index, username: scooby, password: doobie } }YAMLconfigs: elasticSearchUrl: http://localhost:9200 indexName: my_index username: scooby password: doobieElasticsearch 6.2 之前6.2 之前的 Elasticsearch 必须显式设置typeName不能是_docJSON{ elasticSearchUrl: http://localhost:9200, indexName: my_index, typeName: doc, username: scooby, password: doobie }YAMLconfigs: elasticSearchUrl: http://localhost:9200 indexName: my_index typeName: doc username: scooby password: doobieSchema aware 主键示例若要在 Schema aware 模式下把id与a两个字段作为联合主键可参考如下配置出自测试配置 sinkConfig.yaml{ elasticSearchUrl: http://localhost:90902, indexName: myIndex, typeName: doc, username: scooby, password: doobie, primaryFields: id,a }实战从零启动 Elasticsearch sink下面按官方文档 io-elasticsearch-sink.md 的步骤演示完整的本地运行流程。1. 启动单节点 Elasticsearch 集群$ docker run -p 9200:9200 -p 9300:9300 \ -e discovery.typesingle-node \ docker.elastic.co/elasticsearch/elasticsearch:7.13.32. 以 standalone 模式启动本地 Pulsar 服务$ bin/pulsar standalone启动前请确认 NAR 文件存在于connectors/pulsar-io-elastic-search-version.narversion为当前 Pulsar 版本号例如 2.10.x。3. 以 localrun 模式启动连接器使用上文 JSON 配置直接内联$ bin/pulsar-admin sinks localrun \ --archive connectors/pulsar-io-elastic-search-version.nar \ --tenant public \ --namespace default \ --name elasticsearch-test-sink \ --sink-config {elasticSearchUrl:http://localhost:9200,indexName: my_index,username: scooby,password: doobie} \ --inputs elasticsearch_test使用 YAML 配置文件$ bin/pulsar-admin sinks localrun \ --archive connectors/pulsar-io-elastic-search-version.nar \ --tenant public \ --namespace default \ --name elasticsearch-test-sink \ --sink-config-file elasticsearch-sink.yml \ --inputs elasticsearch_test4. 向 topic 发布消息$ bin/pulsar-client produce elasticsearch_test --messages {\a\:1}5. 验证 Elasticsearch 中的文档刷新索引$ curl -s http://localhost:9200/my_index/_refresh检索文档$ curl -s http://localhost:9200/my_index/_search此时可以看到先前发布的消息已成功写入 Elasticsearch返回结果形如{took:2,timed_out:false,_shards:{total:1,successful:1,skipped:0,failed:0},hits:{total:{value:1,relation:eq},max_score:1.0,hits:[{_index:my_index,_type:_doc,_id:FSxemm8BLjG_iC0EeTYJ,_score:1.0,_source:{a:1}}]}}注意示例中 Raw 模式下_id由 Elasticsearch 自动生成FSxemm8BLjG_iC0EeTYJ_source即为消息原始 JSON 内容。源码级要点与测试验证文档_id生成Raw 模式下_id为 null见 ElasticSearchSink.javaSchema aware 模式下_id由 record key 或primaryFields生成。makeIndexRequest()中仅在_id非空时才显式设置文档 ID见 ElasticSearchClient.java。索引自动创建当createIndexIfNeededtrue时首次写入前会通过createIndexIfNeeded()创建索引并用本地indexCache缓存已存在的索引名以避免重复检查创建时按indexNumberOfShards与indexNumberOfReplicas设置分片与副本见 ElasticSearchClient.java。重试与退避连接器内置RandomExponentialRetry随机指数退避重试次数、基础等待时间、最大重试间隔分别由maxRetries、retryBackoffInMs、maxRetryTimeInSec控制相关实现见 RandomExponentialRetry.java 与 RandomExponentialBackoffPolicy.java。不可恢复错误当 bulk 项失败且失败信息包含mapper_parsing_exception、action_request_validation_exception、illegal_argument_exception等“格式错误”关键字时连接器会按malformedDocAction决定 IGNORE/WARN/FAIL见 ElasticSearchClient.java 与 hasIrrecoverableError()进入 FAILED 状态后后续write()会直接抛出IllegalStateException。测试覆盖仓库提供了较完整的测试集验证上述行为例如 ElasticSearchSinkTests.java覆盖schemaEnable、primaryFields单字段/多字段、nullValueAction等场景、ElasticSearchSinkRawDataTests.javaRaw 模式、ElasticSearchClientSslTests.javaTLS、IndexNameFormatterTest.java时间索引测试通过 Testcontainers 拉起真实 Elasticsearch 容器进行端到端验证见 pulsar-io/elastic-search/pom.xml 中的org.testcontainers:elasticsearch依赖。使用建议与注意事项Raw 与 Schema aware 的取舍如果消息本身已经是 JSON 且无需主键去重/更新直接使用默认 Raw 模式即可若需要基于逻辑主键做 INSERT/UPDATE/DELETE如 CDC 同步请开启schemaEnabletrue并配置primaryFields与nullValueAction。索引名合法性省略indexName时连接器会从 topic 名派生索引名请确保 topic 名能通过索引命名转换小写、去掉persistent://tenant/namespace/前缀、≤255 字节、符合命名正则否则启动后写入会失败。Bulk 与背压bulkConcurrentRequests会直接影响内存占用与吞吐需结合实际写入速率调整bulkActions、bulkSizeInMb与bulkFlushIntervalInMs三个触发条件。版本差异Elasticsearch 6.2 之前必须显式设置typeNameindexNumberOfReplicas与connectionIdleTimeoutInMs的默认值请以当前仓库源码ElasticSearchConfig.java为准文档表格中的数值可能存在滞后。凭据与安全配置username后必须同时提供password生产环境建议通过ssl配置块启用 TLS并正确设置 truststore/keystore 路径与密码按需保留hostnameVerificationtrue。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar HBase Sink Connector 完全指南配置、原理与实战写入Apache Pulsar HBase Sink Connector 完全指南配置、原理与实战写入 本文聚焦 Apache Pulsar 官方提供的 HBas消息队列后端流处理Apache Pulsar HBase Sink Connector 实战指南配置详解、部署运行与源码原理Apache Pulsar HBase Sink Connector 实战指南配置详解、部署运行与源码原理 本文围绕 Apache Pulsar 内置的 HB消息队列后端流处理Apache Pulsar RabbitMQ Sink Connector 完全指南配置、部署与源码级原理剖析Apache Pulsar RabbitMQ Sink Connector 完全指南配置、部署与源码级原理剖析 Apache Pulsar 的 RabbitM消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
企业数字化 ERP 产品动态
相关推荐
LDO瞬态响应全解析:电容、ESR与环路设计实战 做电源设计这些年,LDO的瞬态响应一直是个“看着简单、调起来打脸”的指标。表面上看,它无非是负载电流变化时输出电压波动一下,但真正把示波器接上去之后,下冲幅度、恢复时间、振铃相位,每一个都牵扯到环路补偿、输出电… · 2026/9/23 19:43:35
英里换算公里实战项目:搞定3个高频面试题,告别代码报错 英里换算公里实战项目:搞定3个高频面试题,告别代码报错 刚把网上抄来的英里换算代码跑起来,结果控制台直接抛错?别慌,这种“复制粘贴就崩”的情况太常见了。很多工程师卡在单位换算这种看似简单的逻辑上,其实是因为没搞懂背后的精度陷阱和工程化规范。… · 2026/9/23 20:20:16
搞懂头层皮和二层皮的区别,从入门到精通的避坑指南 搞懂头层皮和二层皮的区别,从入门到精通的避坑指南 版本升级后 API 全变了,这是无数开发者在技术进阶路上遇到的第一道鬼门关。很多人卡在“头层皮”的表象逻辑里,以为读懂了文档就能上手,结果一跑代码全是报错。真正的 入门到精通… · 2026/9/23 20:20:10
2019 天天射干 localhost保姆级教程 3步搞定2019天天射干localhost报错速查手册 复制来的代码跑不通不知道怎么调?别慌,这不仅是你的问题,也是无数开发者踩过的坑。针对【2019 天天射干 localhost】这类看似无厘头实则暗藏玄机的报错,我们整理了一份… · 2026/9/23 20:20:03
逾越节速查手册 逾越节源码图解:3步搞懂版本升级API变更原理 逾越节源码图解:3步搞懂版本升级API变更原理 版本升级后 API 全变了,文档翻烂也找不到对应方法,这是无数开发者踩过的坑。别慌,今天用【图解原理】拆解逾越节核心逻辑,从入口到执行链路逐行剖… · 2026/9/23 20:20:03
3招搞定手机怎么下载微信面试难题实战项目解析 3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29