数据库后端流处理【免费下载链接】EventStoreKurrentDB is a database thats engineered for modern software applications and event-driven architectures. Its event-native design simplifies data modeling and preserves data integrity while the integrated streaming engine solves distributed messaging challenges and ensures data consistency.项目地址https://gitcode.com/gh_mirrors/ev/EventStore点击查看免费下载导读本文面向需要在 KurrentDB 与 RabbitMQ 之间搭建事件管道的开发者完整讲解rabbit-mq-sink连接器的创建、参数配置、Broker 确认机制与弹性策略。读完本文你将掌握如何用一条 HTTP 请求把 KurrentDB 流中的事件以指定路由键routing key发布到 RabbitMQ 交换机exchange并理解密码类敏感字段如何被加密保护、以及连接异常时连接器如何依靠 RabbitMQ 自身的重试机制恢复。概览RabbitMQ Sink 的职责RabbitMQ Sink 是 KurrentDB Connectors 插件提供的服务端连接器之一核心职责是将 KurrentDB 事件发布到 RabbitMQ 交换机并使用指定的路由键完成消息路由。它把 RabbitMQ 交换机与队列管理的复杂性封装在连接器内部根据路由键将消息送达与交换机绑定的合适队列从而保证生产环境下消息投递的一致性与可靠性并具备优雅的错误处理与恢复机制。在 KurrentDB 源码中该连接器被注册为RabbitMqSink对应实例类型别名rabbit-mq-sink并挂载了CONNECTORS_RABBITMQ_SINK授权entitlement且要求许可证requiresLicense: true详见 ConnectorCatalogue.cs。从源码结构看与 Elasticsearch、Kafka、MongoDB 等连接器一样它属于需要许可证的商用连接器。前置条件使用 RabbitMQ Sink 连接器之前需要满足以下条件数据保护令牌data protection token必须在 KurrentDB 实例中配置加密令牌用于加密密码等敏感字段可访问的 RabbitMQ 服务端具备可连接的主机host与端口port有效的认证凭据RabbitMQ 用户名与密码。数据保护令牌的具体配置方式参见 Data Protection 文档。以 YAML 配置为例既可以使用令牌文件也可以直接指定令牌Connectors: Enabled: true DataProtection: TokenFile: /path/to/token/fileConnectors: Enabled: true DataProtection: Token: your-secret-token需要特别注意两点来自 features.md 的官方警告令牌一旦设置即为永久性绝不可更改所有连接器共用同一个令牌修改令牌会导致此前加密的数据全部不可访问若同时提供Token与TokenFile系统会优先使用令牌文件并忽略Token设置。底层加密采用 envelope encryption敏感字段先用数据加密密钥DEK加密再由基于令牌的主密钥包裹密文安全存储在 KurrentDB 内部系统流中通过 Surge key vault 管理而非明文落盘。快速开始创建 RabbitMQ Sink 连接器通过 KurrentDB 的管理 API 创建连接器。将{id}替换为自定义的连接器 IDPOST /connectors/{id} Host: localhost:2113 Content-Type: application/json { settings: { instanceTypeName: rabbit-mq-sink, exchange:name: example-exchange, exchange:type: direct, routingKey: my-routing-key, subscription:filter:scope: stream, subscription:filter:filterType: streamId, subscription:filter:expression: example-stream } }创建并启动后每当有事件被追加到example-stream流连接器就会把该记录发送到 RabbitMQ 中指定的交换机进而由路由键路由到相应队列。本例中的订阅过滤器subscription filter含义为以流为作用域scope: stream按流 IDfilterType: streamId匹配example-stream流的事件。订阅过滤的补充说明订阅过滤参数属于所有 Sink 连接器的公共配置详见 Sink Options要点如下subscription:filter:scope取值为stream按流 ID 过滤或record按事件类型过滤subscription:filter:filterType取值为streamId、regex、prefix或jsonPathsubscription:filter:expression过滤表达式若不指定作用域则默认消费$all流并排除系统事件若指定了作用域但表达式为空则消费$all且包含系统事件subscription:initialPosition无既有检查点时消费者从何处开始消费取值为latest或earliest默认latest。连接器生命周期管理创建后可通过管理 API 对连接器进行完整生命周期管理相关端点见 API Reference操作方法与路径说明创建POST /connectors/{id}创建连接器{id}为自定义唯一标识启动POST /connectors/{id}/start启动连接器不带位置参数时从检查点若无则按订阅初始位置开始指定位置启动POST /connectors/{id}/start/{log_position}从指定日志位置开始消费列表GET /connectors支持按state、instanceTypeName、connectorId过滤与分页查看配置GET /connectors/{id}/settings查看当前配置重配置PUT /connectors/{id}/settings无需删除重建即可修改配置运行中修改需重启后生效重置POST /connectors/{id}/reset默认重置到流起点也可reset/{log_position}指定位置停止POST /connectors/{id}/stop停止连接器删除DELETE /connectors/{id}删除连接器重命名PUT /connectors/{id}/rename修改连接器名称配置参数详解除继承自 Sink Options 的公共配置实例类型、订阅、转换、自动提交、日志等外RabbitMQ Sink 还支持以下专属选项Option说明exchange:name必填。交换机名称。exchange:type交换机类型。默认fanoutroutingKey供交换机决定如何将消息路由到合适队列路由键会与交换机所关联队列的绑定规则进行匹配影响消息的投递路径。默认hostRabbitMQ 服务端主机名。默认localhostportRabbitMQ 服务端端口号。默认5672connectionName用于日志与诊断的连接名称。默认连接器 IDvirtualHostRabbitMQ 使用的 VirtualHostvhost。默认/waitForBrokerAck发送操作完成前channel 是否等待 Broker 确认。默认falseauthentication:username认证用户名。默认guestauthentication:password认证密码。默认guestautoDelete交换机不再使用时是否自动删除。默认falsedurable交换机是否为持久化durable。默认true关键概念交换机、路由键与绑定交换机exchangeRabbitMQ 消息进入的入口。exchange:type决定路由语义常见类型包括direct按路由键精确匹配绑定键、fanout广播到所有绑定队列忽略路由键、topic按模式匹配路由键等。默认fanout意味着若不关心路由匹配、希望广播可无需设置路由键。路由键routingKey消息携带的路由信息交换机依据它与队列绑定键的匹配结果决定投递路径。对于direct类型路由键须与某队列绑定键完全一致才会投递。durable 与 autoDeletedurable: true默认保证交换机在 Broker 重启后依然存在autoDelete: false默认保证交换机在无绑定队列后不会被自动删除两者共同保障交换机生命周期稳定可控。弹性Resilience配置RabbitMQ Sink不继承公共的 Resilience configuration而是依赖 RabbitMQ 客户端自身的重试机制并提供以下三个专属弹性参数Option说明resilience:connectionTimeoutMsTCP 建连超时毫秒0表示无限等待。默认60000resilience:handshakeTimeoutMs协议握手超时毫秒。默认10000resilience:requestedHeartbeatMs请求的心跳超时毫秒0表示禁用心跳不应低于 1 秒。默认60000这三项分别约束连接建立的三个阶段TCP 连接建立、AMQP 协议握手、连接保持期间的存活探测。合理调优可避免在网络抖动或 Broker 响应缓慢时出现长时间阻塞或误判断连。Broker Acknowledgment发布确认开关waitForBrokerAck示例中亦可写作WaitForBrokerAck控制连接器是否等待 RabbitMQ Broker 确认后再认为一次发送完成POST /connectors/{id} Host: localhost:2113 Content-Type: application/json { WaitForBrokerAck: true }开启确认默认语义每条消息在发送前均先被 Broker 确认publish 操作才算完成。可靠性优先代价是吞吐量下降——生产者必须等待每个确认才能继续。关闭确认生产者无需等待 Broker 确认即可持续发送消息显著提升吞吐量适合以性能为先的高吞吐场景代价是消息丢失或重复的风险略有上升投递保证相应降低。需要注意的是连接器的检查点checkpoint机制与waitForBrokerAck相互配合连接器将已成功处理的位置周期性写入$connectors/{connector-id}/checkpoints系统流见 CheckpointingBroker 确认与否会影响处理成功的判定时机从而影响重启后的消费起点精确度。连接器整体采用at least once投递保证——事件可能被重复投递但严格按序投递事件x不会被早于其前序事件投递。Headers自动注入的消息头RabbitMQ Sink 会在发送到交换机的每条消息中自动附带一组标准消息头这些头以esdb-前缀标识用于与用户自定义头区分完整列表见 HeadersHeader说明esdb-connector-id处理该记录的唯一连接器 IDesdb-request-id处理请求的唯一标识esdb-request-date请求处理时间戳esdb-record-id记录唯一标识esdb-record-timestamp记录时间戳esdb-record-redacted记录是否被脱敏布尔值esdb-record-stream-id事件来源流 IDesdb-record-stream-revision事件在其流中的修订号esdb-record-log-position事件的全局日志位置esdb-record-schema-subject事件的 schema 主题/名称esdb-record-schema-type事件的 schema 数据格式类型此外若事件带有用户自定义头且连接器允许向 sink 传头这些头会被自动加上esdb-record-headers-前缀默认以$开头的系统头会被排除可通过headers:ignoreSystem: false将其包含到输出中。这些头信息在 RabbitMQ 消费端可用于追踪记录来源、关联请求与做幂等去重。数据保护密码如何被加密RabbitMQ Sink 的authentication:password属于敏感字段会被数据保护框架加密存储。源码中通过RabbitMqSinkConnectorDataProtector将Authentication:Password注册为敏感键sensitive key详见 ConnectorDataProtectors.cspublic class RabbitMqSinkConnectorDataProtector(IDataProtector dataProtector) : ConnectorDataProtectorRabbitMqSinkOptions(dataProtector) { protected override string[] ConfigureSensitiveKeys() [ Authentication:Password ]; }因此即使通过GET /connectors/{id}/settings查看配置密码也不会以明文形式暴露。这也是前置条件中必须首先配置数据保护令牌的根本原因没有令牌敏感字段无法完成加密存储。结合源码的调用关系从源码可以进一步印证本文所述各环节的落点连接器注册ConnectorCatalogue.cs 中RabbitMqSink对应RabbitMqSinkValidator参数校验与RabbitMqSinkConnectorDataProtector敏感字段加密三者配套注册依赖打包KurrentDB.Connectors.csproj 通过Kurrent.Connectors.RabbitMQ包引入 RabbitMQ 客户端能力说明该连接器实质是 KurrentDB 对 RabbitMQ 客户端库的服务端封装插件启用Connector 插件随所有 KurrentDB 二进制预装且默认启用Connectors.Enabled默认true如需禁用可在配置中设置Connectors: Enabled: false见 settings.md。实战小结一个完整的 RabbitMQ Sink 使用流程为配置数据保护令牌 → 创建连接器指定实例类型、交换机、路由键与订阅过滤→ 启动连接器 → 向目标流追加事件 → 事件经订阅过滤后被发布到 RabbitMQ 交换机 → 依据路由键路由到队列。生产环境中建议先配置好数据保护令牌再创建连接器并妥善保管令牌根据路由需求选择exchange:type精确路由用direct广播用fanout高吞吐场景可关闭waitForBrokerAck换取性能可靠性敏感场景保持开启通过网络分区风险较高或 Broker 负载波动大的环境合理调大resilience:connectionTimeoutMs与handshakeTimeoutMs避免误判断连或频繁重连。赞分享数据库后端流处理【免费下载链接】EventStoreKurrentDB is a database thats engineered for modern software applications and event-driven architectures. Its event-native design simplifies data modeling and preserves data integrity while the integrated streaming engine solves distributed messaging challenges and ensures data consistency.项目地址https://gitcode.com/gh_mirrors/ev/EventStore点击查看免费下载相关推荐KurrentDB Kafka Sink 连接器实战指南配置、认证、分区与投递保证KurrentDB Kafka Sink 连接器实战指南配置、认证、分区与投递保证 本指南以 KurrentDBEventStore 项目当前的产品形态中数据库后端流处理KurrentDB HTTP Sink 连接器实战指南配置、认证、动态 URL 与源码原理KurrentDB HTTP Sink 连接器实战指南配置、认证、动态 URL 与源码原理 本指南以 KurrentDB 内置的 HTTP Sink 连接器数据库后端流处理Apache Pulsar Flume Sink 连接器配置与使用指南Apache Pulsar Flume Sink 连接器配置与使用指南 Flume sink 连接器是 Apache Pulsar 的 IO 连接器之一负责将消息队列后端流处理上一篇AList终极指南3步打造你的个人云存储管理中心下一篇10分钟终极指南让AI像人类一样操控Windows界面创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考