ARTICLE DETAIL

资讯详情

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

Apache Pulsar 连接器管理实战:部署、配置、运行与监控 Pulsar IO Connectors

Apache Pulsar 连接器管理实战:部署、配置、运行与监控 Pulsar IO Connectors 消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Apache Pulsar 的Pulsar IO机制让 Pulsar 能够与 Cassandra、Kafka、Aerospike、Kinesis、RabbitMQ 等外部系统双向交换数据Source 连接器负责把外部数据灌入PulsarSink 连接器负责把 Pulsar 数据导出到外部系统。本文以官方文档 io-managing.md 为核心脉络完整讲解如何在 Pulsar 集群中部署内置连接器、编写连接器 YAML 配置、通过pulsar-admin的source/sink命令运行与监控连接器并结合本仓库源码如 Cassandra Sink 的配置实现与 NAR 打包机制给出源码级佐证。读完本文你将掌握 Pulsar IO 连接器的完整生命周期管理技能从内置连接器安装、YAML 配置编写到提交运行、状态监控与升级更新。Pulsar IO 与连接器的基本概念在进入管理操作之前先明确两个核心概念详见 io-overview.mdSource源把数据从外部系统注入Pulsar。常见形态是其他消息系统或firehose 式数据管道 API。Sink汇从 Pulsar消费数据并写入外部系统。常见形态是其他消息系统以及 SQL/NoSQL 数据库。一个关键事实是Pulsar IO 连接器本质上是特化的 Pulsar Functions见 functions-overview.md。因此连接器的管理界面与 Pulsar Functions 高度相似本文后面会看到连接器的监控命令正是复用functions命令族。使用内置连接器安装与自动发现Pulsar 发行版内置了一批经过打包和测试的常用连接器完整清单见 io-connectors.md包括 Aerospike Sink、Cassandra Sink、Kafka Source/Sink、Kinesis Sink、RabbitMQ Source、Twitter Firehose Source 以及基于 Debezium 的 CDC Source 等。自2.1.0-incubating版本起Pulsar 提供独立的内置连接器二进制发行包。安装方式参考 io-quickstart.md 中的Installing Builtin Connectors章节解压连接器发行包后把connectors目录拷贝到 Pulsar 安装根目录# 解压连接器发行包 $ tar xvfz /path/to/apache-pulsar-io-connectors-version-bin.tar.gz # 拷贝 connectors 目录到 Pulsar 根目录 $ cp -r apache-pulsar-io-connectors-version/connectors connectors # 查看可用的内置连接器 NAR 包 $ ls connectors pulsar-io-aerospike-version.nar pulsar-io-cassandra-version.nar pulsar-io-kafka-version.nar pulsar-io-kinesis-version.nar pulsar-io-rabbitmq-version.nar pulsar-io-twitter-version.nar ...安装完成后所有内置连接器都会被 Pulsar broker或 function-worker自动发现无需额外的安装步骤。你可以通过 REST 接口验证连接器是否就绪curl -s http://localhost:8080/admin/v2/functions/connectors返回的 JSON 数组中会列出每个连接器的name、description以及对应的sourceClass/sinkClass例如[{name:aerospike,description:Aerospike database sink,sinkClass:org.apache.pulsar.io.aerospike.AerospikeStringSink}, {name:cassandra,description:Writes data into Cassandra,sinkClass:org.apache.pulsar.io.cassandra.CassandraStringSink}, {name:kafka,description:Kafka source and sink connector,sourceClass:org.apache.pulsar.io.kafka.KafkaStringSource,sinkClass:org.apache.pulsar.io.kafka.KafkaBytesSink}, {name:kinesis,description:Kinesis sink connector,sinkClass:org.apache.pulsar.io.kinesis.KinesisSink}, {name:rabbitmq,description:RabbitMQ source connector,sourceClass:org.apache.pulsar.io.rabbitmq.RabbitMQSource}, {name:twitter,description:Ingest data from Twitter firehose,sourceClass:org.apache.pulsar.io.twitter.TwitterFireHose}]这里返回的name如cassandra、kafka正是后续提交内置连接器时使用的--source-type/--sink-type取值。它的来源是每个连接器 NAR 包内pulsar-io.yaml中声明的name字段详见下文内置类型背后的 NAR 机制。配置连接器YAML 配置文件配置 Pulsar IO 连接器的方式非常直接在运行连接器时提供一个 YAML 配置文件。该文件告诉 PulsarSource/Sink 位于哪里类名或内置类型以及如何把外部系统与 Pulsar topic 连接起来。以官方文档中的 Cassandra Sink 配置为例完整示例见 io-managing.md 与 io-quickstart.md 中的examples/cassandra-sink.ymltenant: public namespace: default name: cassandra-test-sink ... # cassandra 专用配置 configs: roots: localhost:9042 keyspace: pulsar_test_keyspace columnFamily: pulsar_test_table keyname: key columnName: col这份配置表达了三层信息归属连接器运行在哪个tenant、namespace下叫什么名字cassandra-test-sink外部系统位置roots指向 Cassandra 集群地址localhost:9042keyspace与columnFamily指定了数据要写入的键空间与列族/表消息到表的映射keyname和columnName决定 Pulsar 消息如何映射为 Cassandra 表的 key 与列。这些configs字段并非随手定义而是由连接器实现类的配置类严格约束的。以本仓库为例Cassandra Sink 的配置类在 CassandraSinkConfig.javaFieldDoc(required true, help A comma-separated list of cassandra hosts to connect to) private String roots; // Cassandra 主机列表逗号分隔 FieldDoc(required true, help The key space used for writing pulsar messages to) private String keyspace; // 写入 Pulsar 消息所用的键空间 FieldDoc(required true, help The key name of the cassandra column family) private String keyname; // Cassandra 列族中的 key 名 FieldDoc(required true, help The cassandra column family name) private String columnFamily; // Cassandra 列族表名 FieldDoc(required true, help The column name of the cassandra column family) private String columnName; // Cassandra 列族中的列名从源码可以看到这 5 个字段全部标记为required true即缺少任何一个都会导致配置校验失败同时该类提供了load(String yamlFile)与load(MapString, Object map)两种加载方式分别支持从 YAML 文件路径和键值 Map 构造配置对应命令行中的--sink-config-file与--sink-config两种传参方式。每个内置连接器的专用配置项可能不同使用时请查阅对应连接器的独立文档内置连接器清单与文档索引见 io-overview.md 的 Working with connectors 表格。运行连接器连接器的运行管理统一通过pulsar-adminCLI 工具详见 reference-pulsar-admin.md中的source和sink两个命令族完成。每个命令族都包含create、update、delete、localrun、available-sources/available-sinks等子命令。运行 Source提交到集群运行的自定义 Source命令形式如下./bin/pulsar-admin source create \ --classname classname \ --archive jar-location \ --tenant tenant \ --namespace namespace \ --name source-name \ --destination-topic-name output-topic官方文档给出的完整示例提交 Twitter Firehose Sourcebin/pulsar-admin source create \ --classname org.apache.pulsar.io.twitter.TwitterFireHose \ --archive ~/application.jar \ --tenant test \ --namespace ns1 \ --name twitter-source \ --destination-topic-name twitter_data其中--classname指向 Twitter Source 的实现类org.apache.pulsar.io.twitter.TwitterFireHose源码位于 pulsar-io/twitter 模块--archive指向打包好的 JAR/NAR--destination-topic-name指定 Source 产出的数据写入哪个 Pulsar topic。在本地以独立进程运行不提交到集群把create换成localrun即可bin/pulsar-admin source localrun \ --classname org.apache.pulsar.io.twitter.TwitterFireHose \ --archive ~/application.jar \ --tenant test \ --namespace ns1 \ --name twitter-source \ --destination-topic-name twitter_data提交内置 Source时无需指定--classname和--archive只需通过--source-type声明内置类型./bin/pulsar-admin source create \ --tenant tenant \ --namespace namespace \ --name source-name \ --destination-topic-name input-topics \ --source-type source-type例如提交一个 Kafka Source把 Kafka 数据导入 Pulsar topic./bin/pulsar-admin source create \ --tenant test-tenant \ --namespace test-namespace \ --name test-kafka-source \ --destination-topic-name pulsar_sink_topic \ --source-type kafkasource create还支持以下常用参数完整选项表见 reference-pulsar-admin.md 的source章节Flag说明--deserialization-classnameSource 的 SerDe 类名--schema-type内置 schema如avro、json或自定义 Schema 类名用于编码 Source 产出的消息--source-config以键值对方式传入 Source 配置--source-config-fileYAML 配置文件路径--parallelismSource 实例数并行度--processing-guarantees处理语义可选ATLEAST_ONCE、ATMOST_ONCE、EFFECTIVELY_ONCE运行 Sink提交到集群运行的自定义 Sink./bin/pulsar-admin sink create \ --classname classname \ --archive jar-location \ --tenant test \ --namespace namespace \ --name sink-name \ --inputs input-topics官方文档示例提交 Cassandra Sink./bin/pulsar-admin sink create \ --classname org.apache.pulsar.io.cassandra \ --archive ~/application.jar \ --tenant test \ --namespace ns1 \ --name cassandra-sink \ --inputs test_topic本地运行同样是换成localrun./bin/pulsar-admin sink localrun \ --classname org.apache.pulsar.io.cassandra \ --archive ~/application.jar \ --tenant test \ --namespace ns1 \ --name cassandra-sink \ --inputs test_topic提交内置 Sink时使用--sink-type./bin/pulsar-admin sink create \ --tenant tenant \ --namespace namespace \ --name sink-name \ --inputs input-topics \ --sink-type sink-type示例提交 Cassandra Sink./bin/pulsar-admin sink create \ --tenant test-tenant \ --namespace test-namespace \ --name test-cassandra-sink \ --inputs pulsar_input_topic \ --sink-type cassandra注意内置连接器的sink-type参数取值由连接器 NAR 包中pulsar-io.yaml文件里name参数的设置决定同一规则也适用于source-type。sink create的常用参数完整选项表见 reference-pulsar-admin.md 的sink章节Flag说明--inputsSink 的输入 topic多个 topic 用逗号分隔--topics-pattern按正则模式消费某 namespace 下的多个 topic--sink-config/--sink-config-file键值对方式 / YAML 文件方式传入 Sink 配置--custom-serde-inputs输入 topic 到 SerDe 类名的映射JSON 字符串--custom-schema-inputs输入 topic 到 Schema 类型或类名的映射JSON 字符串--parallelismSink 实例数并行度--processing-guarantees处理语义可选ATLEAST_ONCE、ATMOST_ONCE、EFFECTIVELY_ONCE--auto-ack是否由 Functions 框架统一管理 ack--timeout-ms消息超时时间毫秒查看可用内置连接器两个命令族都提供了枚举内置连接器的子命令可用于确认集群中可用的内置类型bin/pulsar-admin sink available-sinks bin/pulsar-admin source available-sources监控连接器由于 Pulsar IO 连接器以 Pulsar Functions 的形式运行因此可以直接复用pulsar-admin的functions命令族进行监控functions命令的完整说明同样见 reference-pulsar-admin.md。获取连接器元数据bin/pulsar-admin functions get \ --tenant tenant \ --namespace namespace \ --name connector-name以官方文档的 Cassandra Sink 为例返回的 JSON 会展示连接器的归属、类名、topic 绑定与配置{ tenant: public, namespace: default, name: cassandra-test-sink, className: org.apache.pulsar.functions.api.utils.IdentityFunction, autoAck: true, parallelism: 1, source: { topicsToSerDeClassName: { test_cassandra: } }, sink: { configs: {\roots\:\cassandra\,\keyspace\:\pulsar_test_keyspace\,\columnFamily\:\pulsar_test_table\,\keyname\:\key\,\columnName\:\col\}, builtin: cassandra }, resources: {} }注意输出中sink.builtin字段值为cassandra印证了--sink-type cassandra的提交结果sink.configs中保存的就是 YAML 中configs段的序列化内容。获取连接器运行状态bin/pulsar-admin functions getstatus \ --tenant tenant \ --namespace namespace \ --name connector-name刚提交尚未处理消息时状态输出大致如下{ functionStatusList: [ { running: true, instanceId: 0, metrics: { metrics: { __total_processed__: {}, __total_successfully_processed__: {}, __total_system_exceptions__: {}, __total_user_exceptions__: {}, __total_serialization_exceptions__: {}, __avg_latency_ms__: {} } }, workerId: c-standalone-fw-localhost-6750 } ] }当向输入 topictest_cassandra生产 10 条消息后bin/pulsar-client produce -m key-$i -n 1 test_cassandra再次getstatus可以看到numProcessed、numSuccessfullyProcessed增长到 11多出的 1 条来自 quickstart 示例中先前的验证消息__total_processed__等指标开始有值——这是判断连接器是否正常消费并写出的直接依据。基于指标的深度监控除了上述命令连接器运行健康度还可以从指标层面持续观测详见 io-develop.md 的 Monitor 章节Pulsar 自带指标Java 连接器暴露的指标可被 Prometheus 采集、在 Grafana 中查看监控部署指南见 deploy-monitoring.md自定义指标Java 连接器可在open阶段通过SinkContext.recordMetric(name, value)上报自定义指标function-worker 会自动将其收集到 Prometheus。public class TestMetricSink implements SinkString { Override public void open(MapString, Object config, SinkContext sinkContext) throws Exception { sinkContext.recordMetric(foo, 1); } Override public void write(RecordString record) throws Exception { } Override public void close() throws Exception { } }更新、删除与升级连接器管理连接器的完整生命周期还包含更新、删除与升级。source/sink命令族均提供了对应子命令选项与create基本一致更新升级使用source update/sink update可以更新已提交的连接器——例如更换--archive指向新版本 JAR/NAR、调整--parallelism、修改--source-config-file/--sink-config-file或切换--source-type/--sink-type这是官方文档Upgrade a connector目标的落地方式bin/pulsar-admin sink update \ --tenant public \ --namespace default \ --name cassandra-test-sink \ --sink-type cassandra \ --sink-config-file examples/cassandra-sink-v2.yml \ --inputs test_cassandra删除使用sink delete/source delete停止并移除连接器quickstart 示例中的收尾操作bin/pulsar-admin sink delete \ --tenant public \ --namespace default \ --name cassandra-test-sink升级自定义连接器时需把新版本打包为 JAR/NAR 并替换--archive路径升级内置连接器时则需替换connectors目录下对应的 NAR 文件后重新update提交。具体更新选项请参考 reference-pulsar-admin.md 中source update/sink update的选项表。内置类型背后的 NAR 机制pulsar-io.yaml理解--source-type/--sink-type的来源有助于自定义连接器的部署。Pulsar 采用NARNiFi Archive打包机制详见 io-develop.md 的 Package 章节所有内置连接器均以 NAR 形式发布。每个 NAR 包内必须包含resources/META-INF/services/pulsar-io.yaml描述文件其结构如下name: connector name description: connector description sourceClass: fully qualified class name (only if source connector) sinkClass: fully qualified class name (only if sink connector)以本仓库的 Cassandra Sink 为例pulsar-io/cassandra/src/main/resources/META-INF/services/pulsar-io.yaml 声明name: cassandra description: Writes data into Cassandra sinkClass: org.apache.pulsar.io.cassandra.CassandraStringSink sinkConfigClass: org.apache.pulsar.io.cassandra.CassandraSinkConfig这里name: cassandra正是提交内置 Sink 时使用的--sink-type cassandra取值也是 REST 接口/admin/v2/functions/connectors返回列表中的name。同时sinkConfigClass指明了该连接器的配置类 CassandraSinkConfig.javaPulsar IO 运行时据此解析 YAML 中configs段的字段。因此部署自定义连接器时只需正确编写pulsar-io.yaml并用 NAR 方式打包就能获得与内置连接器一致的--source-type/--sink-type提交体验。小结Pulsar IO 连接器的管理可以归纳为四步闭环部署安装内置连接器发行包到 Pulsar 根目录的connectors目录broker/function-worker 自动发现配置编写 YAML 配置文件configs段字段由连接器实现类的配置类如CassandraSinkConfig严格定义运行用pulsar-admin source/sink命令族的create提交集群、localrun本地进程提交连接器内置连接器只需--source-type/--sink-type监控与迭代用functions get/getstatus查看元数据与运行状态用source/sink update升级、source/sink delete下线。由于连接器即 Functions你还可以进一步参考 functions-deploying.md、functions-guarantees.md 等文档深入理解处理语义与部署细节自定义连接器的开发、Schema 处理与 NAR 打包规范则见 io-develop.md。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar IO 连接器管理实战内置连接器部署、配置、运行与监控Apache Pulsar IO 连接器管理实战内置连接器部署、配置、运行与监控 本文以 Apache Pulsar IOPulsar IO framewo消息队列后端流处理Apache Pulsar IO 连接器管理实战内置连接器的部署、运行、监控与升级Pulsar 2.3.0Apache Pulsar IO 连接器管理实战内置连接器的部署、运行、监控与升级Pulsar 2.3.0 本篇指南聚焦 Apache Pulsar IO消息队列后端流处理Apache Pulsar 连接器管理实战内置连接器部署、配置、运行与监控Apache Pulsar 连接器管理实战内置连接器部署、配置、运行与监控 本篇技术指南聚焦 Apache Pulsar 的 IO Connectors连接消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表