
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Apache Pulsar 的 IO 连接器Pulsar IO Connectors用于打通 Pulsar 与外部数据系统数据库、其他消息系统等是构建数据管道的关键能力。本文以官方 2.3.2 版文档为主体结合本仓库源码系统讲解 Source / Sink 概念、三种处理保证Processing Guarantees语义及其设置与更新方法并延伸介绍连接器的部署、管理与监控实操。读完本文你将能够理解 Pulsar IO 的架构模型掌握通过pulsar-adminCLI 创建、更新、删除、查看 Source 与 Sink 的完整操作方式。Pulsar IO 连接器消息系统与外部世界之间的桥梁消息系统Messaging System只有能轻松与数据库、其他消息系统等外部系统对接时才最能发挥其价值。Pulsar IO 连接器正是为此而生它让你能够轻松地创建create、部署deploy和管理manage与外部系统交互的连接器例如 Apache Cassandra、Aerospike 以及许多其他系统。从整体上看Pulsar IO 连接器在功能形态上属于 Pulsar Functions 体系连接器与函数Functions都是实例instance的组成部分并且都运行在 Functions Worker 上。当你通过 Connector Admin CLI即pulsar-admin sources/pulsar-admin sinks或 Functions Admin CLI 管理一个 Source、Sink 或 Function 时会在某个 Worker 上启动一个实例。关于 Functions Worker 的更多细节可参考 Functions worker。Source 与 Sink连接器的两种类型Pulsar IO 连接器分为两种类型Source数据源和Sink数据汇。下面的示意图展示了 Source、Pulsar 与 Sink 之间的关系Pulsar IO 连接器Source 与 Sink架构示意图)Source将外部数据流入 PulsarSource 的职责是将外部系统中的数据馈入 Pulsarfeed data from external systems into Pulsar。常见的 Source 包括其他消息系统以及 firehose 风格的数据管道 API例如 Twitter Firehose。一个典型的例子是 Kafka Source它从 Kafka 主题消费消息再写入 Pulsar 主题。Pulsar 内置 Source 连接器的完整清单参见 source connector。从本仓库的目录结构看内置连接器分布在 pulsar-io 目录下例如 Kafka Source、RabbitMQ Source、Twitter Firehose 等每个连接器都是一个独立的 Maven 模块。Sink将 Pulsar 数据流出到外部系统Sink 的职责是将 Pulsar 中的数据馈入外部系统feed data from Pulsar into external systems。常见的 Sink 包括其他消息系统以及 SQL、NoSQL 数据库例如 Cassandra Sink、ElasticSearch Sink、HBase Sink、Redis Sink、MongoDB Sink、InfluxDB Sink 等。Pulsar 内置 Sink 连接器的完整清单参见 sink connector。内置连接器清单版本 2.3.2根据 Builtin Connectors 文档Pulsar 发行版内置了经过打包和测试的一组常用连接器覆盖大多数常见数据系统Aerospike Sink ConnectorCassandra Sink ConnectorKafka Sink Connector / Kafka Source ConnectorKinesis Sink ConnectorRabbitMQ Source Connector / RabbitMQ Sink ConnectorTwitter Firehose Source ConnectorCDC Source Connector基于 DebeziumNetty Source ConnectorHBase Sink ConnectorElasticSearch Sink ConnectorFile Source ConnectorHDFS Sink ConnectorMongoDB Sink ConnectorRedis Sink ConnectorSolr Sink ConnectorInfluxDB Sink Connector处理保证Processing Guarantees连接器的消息语义处理保证Processing Guarantees用于处理向 Pulsar 主题写入消息时发生的错误。需要特别注意的是Pulsar 连接器Source 与 Sink与 Pulsar Functions 使用相同的处理保证语义具体如下表所示交付语义Delivery semantic描述Descriptionat-most-once至多一次发送给连接器的每条消息要么被处理一次要么不被处理。at-least-once至少一次发送给连接器的每条消息可能被处理一次也可能被处理多次。effectively-once恰好一次 / 有效一次发送给连接器的每条消息对应一个输出。这三种语义同样出现在 Pulsar Functions 的处理保证文档 functions-guarantees.md 中二者保持一致。处理保证的适用范围与边界需要强调的是连接器的处理保证不仅依赖于 Pulsar 自身的保证还与外部系统的实现即 Source 和 Sink 的具体实现密切相关。SourcePulsar 保证向 Pulsar 主题写入消息时遵守处理保证。这一步完全在 Pulsar 的控制范围之内。Sink处理保证依赖于 Sink 的实现。如果 Sink 实现没有以幂等idempotent的方式处理重试那么该 Sink 就无法兑现处理保证。从源码实现上看这一语义在 Function.proto 中被定义为枚举类型ProcessingGuarantees包含三个取值enum ProcessingGuarantees { ATLEAST_ONCE 0; // [default value] ATMOST_ONCE 1; EFFECTIVELY_ONCE 2; }其中ATLEAST_ONCE是枚举的默认值default value与文档中“未指定时默认ATLEAST_ONCE”的描述一一对应。该枚举同时被FunctionDetails消息引用字段processingGuarantees编号 6是整个函数/连接器运行时语义的基础。Source 配置在 SourceConfigUtils.java 中通过convertProcessingGuarantee把用户配置转换为 proto 枚举转换逻辑定义在 FunctionCommon.java。在运行时处理保证语义直接决定了 Source 消费消息后的确认行为。以 PulsarSource.java 为例当处理保证为EFFECTIVELY_ONCE时Source 使用acknowledgeCumulativeAsync累积确认来确认消息否则使用acknowledgeAsync单条确认。而在消息处理失败时若为EFFECTIVELY_ONCE会直接抛出运行时异常其他语义下则对消息执行negativeAcknowledge负向确认触发重新投递——这正是at-least-once语义下“消息可能被处理多次”的实现基础。设置处理保证Set创建连接器时可以通过以下语义值设置处理保证ATLEAST_ONCEATMOST_ONCEEFFECTIVELY_ONCE如果创建连接器时未指定--processing-guarantees参数默认语义为ATLEAST_ONCE。以下以Admin CLI为例展示设置方法REST API 与 Java Admin API 的对应用法可参考 io-use.md#create。创建Source并设置处理保证$ bin/pulsar-admin sources create \ --processing-guarantees ATMOST_ONCE \ # 其他 source 配置项pulsar-admin sources create的完整参数列表参见 reference-connector-admin.md#create。创建Sink并设置处理保证$ bin/pulsar-admin sinks create \ --processing-guarantees EFFECTIVELY_ONCE \ # 其他 sink 配置项pulsar-admin sinks create的完整参数列表参见 reference-connector-admin.md#create-1。更新处理保证Update连接器创建之后可以通过以下命令更新其处理保证取值同样为ATLEAST_ONCEATMOST_ONCEEFFECTIVELY_ONCE更新Source的处理保证$ bin/pulsar-admin sources update \ --processing-guarantees EFFECTIVELY_ONCE \ # 其他 source 配置项pulsar-admin sources update的完整参数列表参见 reference-connector-admin.md#update。更新Sink的处理保证$ bin/pulsar-admin sinks update \ --processing-guarantees ATMOST_ONCE \ # 其他 sink 配置项pulsar-admin sinks update的完整参数列表参见 reference-connector-admin.md#update-1。与连接器协作创建、运行与管理你可以通过 Connector Admin CLI使用 sources 和 sinks 子命令来管理 Pulsar 连接器例如创建、更新、启动、停止、重启、重新加载、删除以及其他操作。使用内置连接器Pulsar 内置了多个 内置连接器用于在常用系统如数据库、消息系统之间移动数据。使用内置连接器非常简单按照 getting-started-standalone.md 中安装内置连接器的说明 完成安装后所有内置连接器都会被 Pulsar Broker或 Functions Worker自动发现无需额外的安装步骤。从 2.3.0 版本开始Pulsar 将全部内置连接器作为独立的 NAR 归档发布。启用这些内置连接器需要先从下载页面下载连接器 NAR 归档然后将其放入解压后的 Pulsar 发行版目录下的connectors目录# 解压 Pulsar tarball 并将连接器 NAR 归档复制到 connectors 目录 $ tar xvfz /path/to/apache-pulsar-version-bin.tar.gz $ cd apache-pulsar-version $ mkdir connectors $ cp -r /path/to/downloaded/connectors/*.nar ./connectors $ 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 ...也可以直接使用自带全部内置连接器的 Docker 镜像apachepulsar/pulsar-all:version。每个连接器模块的pulsar-io.yaml文件位于META-INF/services/下声明了连接器的类型信息。例如 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这也解释了文档中的一条注意事项内置连接器的sink-type/source-type参数值由pulsar-io.yaml中的name参数决定。例如 Cassandra Sink 的--sink-type cassandra即来源于此文件中的name: cassandra。配置连接器YAML 配置文件配置 Pulsar IO 连接器非常简单运行连接器时提供一个 YAML 配置文件即可。该 YAML 配置文件告诉 Pulsar 去哪里定位 Source 和 Sink以及如何将 Source / Sink 与 Pulsar 主题连接起来。下面是一个 Cassandra Sink 的 YAML 配置示例tenant: public namespace: default name: cassandra-test-sink ... # cassandra 专用配置 configs: roots: localhost:9042 keyspace: pulsar_test_keyspace columnFamily: pulsar_test_table keyname: key columnName: col这个示例告诉 Pulsar 要连接哪个 Cassandra 集群roots、在 Cassandra 中使用哪个keyspace和columnFamily来收集数据以及如何将 Pulsar 消息映射为 Cassandra 表的 key 和列keyname、columnName。各连接器的详细配置请查阅 io-overview.md 中连接器清单 对应的各连接器文档。运行 Source 连接器可以使用如下形式的命令将一个 Source 提交到现有 Pulsar 集群中运行$ ./bin/pulsar-admin sources create --classname classname --archive jar-location --tenant tenant --namespace namespace --name source-name --destination-topic-name output-topic示例bin/pulsar-admin sources create --classname org.apache.pulsar.io.twitter.TwitterFireHose --archive ~/application.jar --tenant test --namespace ns1 --name twitter-source --destination-topic-name twitter_data除了提交到集群运行也可以将 Source 作为本地机器上的独立进程运行bin/pulsar-admin sources 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 sources create \ --tenant tenant \ --namespace namespace \ --name source-name \ --destination-topic-name input-topics \ --source-type source-type以提交一个 Kafka Source 为例./bin/pulsar-admin sources create \ --tenant test-tenant \ --namespace test-namespace \ --name test-kafka-source \ --destination-topic-name pulsar_sink_topic \ --source-type kafka运行 Sink 连接器可以使用如下形式的命令将一个 Sink 提交到现有 Pulsar 集群中运行./bin/pulsar-admin sinks create --classname classname --archive jar-location --tenant test --namespace namespace --name sink-name --inputs input-topics示例./bin/pulsar-admin sinks create --classname org.apache.pulsar.io.cassandra --archive ~/application.jar --tenant test --namespace ns1 --name cassandra-sink --inputs test_topic同样也可以将 Sink 作为本地进程运行./bin/pulsar-admin sinks localrun --classname org.apache.pulsar.io.cassandra --archive ~/application.jar --tenant test --namespace ns1 --name cassandra-sink --inputs test_topic如果提交的是内置 Sink无需指定--classname和--archive只需通过--sink-type指定类型注意内置连接器的sink-type参数值由pulsar-io.yaml文件中name参数的设置决定。./bin/pulsar-admin sinks create \ --tenant tenant \ --namespace namespace \ --name sink-name \ --inputs input-topics \ --sink-type sink-type以提交一个 Cassandra Sink 为例./bin/pulsar-admin sinks create \ --tenant test-tenant \ --namespace test-namespace \ --name test-cassandra-sink \ --inputs pulsar_input_topic \ --sink-type cassandra监控连接器查看元数据与运行状态由于 Pulsar IO 连接器是以 Pulsar Functions 的形式运行的因此可以使用 pulsar-admin CLI 工具中的functions命令来监控它们。获取连接器元数据Metadatabin/pulsar-admin functions get \ --tenant tenant \ --namespace namespace \ --name connector-name获取连接器运行状态Statusbin/pulsar-admin functions getstatus \ --tenant tenant \ --namespace namespace \ --name connector-name此外pulsar-admin sources与pulsar-admin sinks子命令也分别提供了status、get等操作。例如在 Cassandra Sink 示例中sink status的返回结果包含numInstances实例总数、numRunning运行中实例数以及每个实例的numReadFromPulsar已从 Pulsar 读取的消息数、numWrittenToSink已写入 Sink 的消息数等统计字段可用于判断连接器的工作状态。实战连接 Pulsar 与 Apache Cassandra下面通过一个端到端示例演示如何使用内置 Cassandra Sink 将 Pulsar 主题中的数据写入 Cassandra全程无需编写任何代码。完整步骤参考 io-quickstart.md。前提以下操作假设 Pulsar 以 standalone 模式 运行并且所有命令都在 Pulsar 二进制发行版的根目录下执行。这些命令同样适用于多节点 Pulsar 集群。启动 Pulsar 服务bin/pulsar standalone所有 Pulsar 组件将按顺序启动。可以通过下面的端点检查服务是否正常检查 Pulsar 二进制协议端口默认 6650telnet localhost 6650检查 Pulsar Functions 集群curl -s http://localhost:8080/admin/v2/worker/cluster确认publictenant 与defaultnamespace 存在curl -s http://localhost:8080/admin/v2/namespaces/public确认所有内置连接器均可用curl -s http://localhost:8080/admin/v2/functions/connectors其中第 4 步的示例输出JSON 数组列出了每个内置连接器的名称、描述以及对应的 Source/Sink 实现类例如[{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}]如果启动过程中出错可以在运行pulsar standalone的终端看到异常也可以查看 Pulsar 目录下logs目录中的日志。启动单节点 Cassandra 集群使用 Docker 启动一个单节点 Cassandra 集群docker run -d --rm --namecassandra -p 9042:9042 cassandra启动后依次验证docker ps确认进程在运行docker logs cassandra检查日志docker exec cassandra nodetool status检查集群状态输出中节点应显示为UN即 Up/Normal。使用cqlsh连接并创建 keyspace 与表$ docker exec -ti cassandra cqlsh localhost Connected to Test Cluster at localhost:9042. [cqlsh 5.0.1 | Cassandra 3.11.2 | CQL spec 3.4.4 | Native protocol v4]cqlsh CREATE KEYSPACE pulsar_test_keyspace WITH replication {class:SimpleStrategy, replication_factor:1}; cqlsh USE pulsar_test_keyspace; cqlsh:pulsar_test_keyspace CREATE TABLE pulsar_test_table (key text PRIMARY KEY, col text);配置并提交 Cassandra Sink创建配置文件examples/cassandra-sink.ymlconfigs: roots: localhost:9042 keyspace: pulsar_test_keyspace columnFamily: pulsar_test_table keyname: key columnName: col提交 Cassandra Sinksink-type取自 pulsar-io.yaml 中的namebin/pulsar-admin sink create \ --tenant public \ --namespace default \ --name cassandra-test-sink \ --sink-type cassandra \ --sink-config-file examples/cassandra-sink.yml \ --inputs test_cassandra命令执行后Pulsar 会创建一个名为cassandra-test-sink的 Sink 连接器它将以 Pulsar Function 的形式运行并把主题test_cassandra中产生的消息写入 Cassandra 表pulsar_test_table。查看 Sink 信息与状态bin/pulsar-admin sink get \ --tenant public \ --namespace default \ --name cassandra-test-sink示例输出{ tenant: public, namespace: default, name: cassandra-test-sink, className: org.apache.pulsar.io.cassandra.CassandraStringSink, inputSpecs: { test_cassandra: { isRegexPattern: false } }, configs: { roots: localhost:9042, keyspace: pulsar_test_keyspace, columnFamily: pulsar_test_table, keyname: key, columnName: col }, parallelism: 1, processingGuarantees: ATLEAST_ONCE, retainOrdering: false, autoAck: true, archive: builtin://cassandra }可以看到未显式指定处理保证时processingGuarantees默认即为ATLEAST_ONCE与本文前面介绍的默认语义完全一致archive: builtin://cassandra说明这是一个内置连接器。检查运行状态bin/pulsar-admin sink status \ --tenant public \ --namespace default \ --name cassandra-test-sink生产消息并验证数据落库向 Sink 的输入主题test_cassandra生产 10 条消息for i in {0..9}; do bin/pulsar-client produce -m key-$i -n 1 test_cassandra; done再次查看 Sink 状态可以看到numReadFromPulsar与numWrittenToSink均为 10{ numInstances : 1, numRunning : 1, instances : [ { instanceId : 0, status : { running : true, error : , numRestarts : 0, numReadFromPulsar : 10, numSystemExceptions : 0, latestSystemExceptions : [ ], numSinkExceptions : 0, latestSinkExceptions : [ ], numWrittenToSink : 10, lastReceivedTime : 1551685489136, workerId : c-standalone-fw-localhost-8080 } } ] }最后在 Cassandra 中验证数据docker exec -ti cassandra cqlsh localhostcqlsh use pulsar_test_keyspace; cqlsh:pulsar_test_keyspace select * from pulsar_test_table; key | col ---------------- key-5 | key-5 key-0 | key-0 ...删除 Sinkbin/pulsar-admin sink delete \ --tenant public \ --namespace default \ --name cassandra-test-sink从源码理解连接器的实现基础运行时模型连接器本质上以 Function 的形态运行。定义连接器运行细节的核心模型位于 Function.proto其中SourceSpec与SinkSpec消息分别承载 Source / Sink 的类名、配置JSON 格式的configs、内置连接器标识builtin、订阅类型等字段FunctionDetails统一管理租户、命名空间、名称、并行度parallelism与处理保证processingGuarantees。处理保证的落点Source 的配置转换在 SourceConfigUtils.java 中完成Sink 对应 SinkConfigUtils.java两者通过 FunctionCommon.java 中的convertProcessingGuarantee在用户配置与 proto 枚举之间互转。运行时PulsarSource.java 依据处理保证选择累积确认或单条确认PulsarSink.java 负责将消息写入外部系统。测试验证仓库中 SourceConfigUtilsTest.java 与 SinkConfigUtilsTest.java 覆盖了配置转换逻辑PulsarSourceTest.java 与 PulsarSinkTest.java 覆盖了运行时行为可作为深入阅读与二次开发的入口。总结Pulsar IO 连接器通过 Source / Sink 两种抽象打通了 Pulsar 与外部数据系统之间的双向数据通道其处理保证at-most-once、at-least-once、effectively-once与 Pulsar Functions 保持统一默认语义为ATLEAST_ONCE可通过pulsar-admin sources/sinks create与update命令设置或更新。内置连接器以 NAR 归档形式发布被 Functions Worker 自动发现连接器以 Function 实例的形态运行在 Worker 上既可用sources/sinks子命令管理也可用functions子命令监控。结合本仓库源码读者可以进一步深入理解处理保证在消费确认、失败重试等环节的具体落点从而在生产环境中正确选择与配置连接器。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar IO Connectors 全景指南Source 与 Sink 架构、处理语义与 Connector Admin 实操Apache Pulsar IO Connectors 全景指南Source 与 Sink 架构、处理语义与 Connector Admin 实操 Apach消息队列后端流处理Apache Pulsar IO 连接器详解Source、Sink 架构与处理语义配置实战Apache Pulsar IO 连接器详解Source、Sink 架构与处理语义配置实战 Pulsar IO 是 Apache Pulsar 内置的流数据接消息队列后端流处理Apache Pulsar IO 连接器概览Source/Sink 架构与处理语义Processing Guarantees深度解析Apache Pulsar IO 连接器概览Source/Sink 架构与处理语义Processing Guarantees深度解析 本文是 Apache消息队列后端流处理上一篇YouTransfer文件传输核心功能详解上传、下载、打包全流程下一篇【免费下载】 QuickDraw: 人工智能简笔画识别项目创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考