ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

EventBridge Pipes 完全仿真指南:Floci 中的管道配置、Kafka 桥接与 Enrichment 实战

EventBridge Pipes 完全仿真指南:Floci 中的管道配置、Kafka 桥接与 Enrichment 实战 EventBridge Pipes 完全仿真指南Floci 中的管道配置、Kafka 桥接与 Enrichment 实战【免费下载链接】flociLight, fluffy, and always free - The AWS Local Emulator alternative项目地址: https://gitcode.com/gh_mirrors/fl/flociFlociLight, fluffy, and always free — AWS 本地模拟器替代品提供了一套与 AWS EventBridge Pipes 线协议兼容的 REST-JSON API让你在本地的 SQS、Kinesis、DynamoDB Streams、Kafka 等源与 Lambda、SQS、SNS、Kinesis、Step Functions 等目标之间搭建可过滤、可富化Enrichment的事件管道。本指南以 docs/services/pipes.md 为主线结合 PipesController.java、PipesService.java、PipesPoller.java 等源码帮助你掌握创建、描述、更新、启停与删除 Pipe 的全部操作理解 Kafka 源如何借助 Karapace 边车工作以及 Enrichment、Filter、ParallelizationFactor 等参数的精确语义与当前限制。概览协议、端点与支持的动作EventBridge Pipes 在 Floci 中通过REST-JSON协议暴露与需要X-Amz-Target的 JSON 1.1 协议不同固定端点为POST http://localhost:4566/从 PipesController.java 可以看到控制器使用标准 HTTP 动词 JSON 请求体并统一以application/json生产/消费。端点全部挂载在/v1/pipes下与 AWS CLI 的pipes子命令一一对应。支持的动作如下表ActionHTTP 路由说明CreatePipePOST /v1/pipes/{name}以 source、target 与可选 enrichment 创建新的管道DescribePipeGET /v1/pipes/{name}获取管道详情含状态与配置UpdatePipePUT /v1/pipes/{name}更新管道配置source、target、role、enrichment、desired state 等DeletePipeDELETE /v1/pipes/{name}删除管道ListPipesGET /v1/pipes列出管道支持按状态与前缀过滤StartPipePOST /v1/pipes/{name}/start启动已停止的管道StopPipePOST /v1/pipes/{name}/stop停止运行中的管道请求体字段采用 AWS 线格式的大驼峰命名Source、Target、RoleArn、Description、Enrichment、DesiredState、SourceParameters、TargetParameters、EnrichmentParameters、Tags。控制器对这些字段逐个解析并透传给 PipesService。快速上手SQS 到 Lambda 的管道准备环境先设置 AWS 端点环境变量指向本地运行的 Flociexport AWS_ENDPOINT_URLhttp://localhost:4566创建管道SQS → Lambda# Create a pipe (SQS to Lambda) aws pipes create-pipe \ --name my-pipe \ --source arn:aws:sqs:us-east-1:000000000000:source-queue \ --target arn:aws:lambda:us-east-1:000000000000:function:my-function \ --role-arn arn:aws:iam::000000000000:role/pipe-role \ --endpoint-url $AWS_ENDPOINT_URLCreatePipe有四个必填字段PipesService.createPipe 中会严格校验Name管道名称不能为空Source源 ARN不能为空Target目标 ARN不能为空RoleArnIAM 角色 ARN不能为空。任一缺失都会返回ValidationExceptionHTTP 400。若同名管道已存在按region::name作为存储键则返回ConflictExceptionHTTP 409。创建时若未指定DesiredState默认值为RUNNING管道创建后会立即被 PipesPoller 调度轮询。描述 / 列表 / 启停 / 更新 / 删除# Describe a pipe aws pipes describe-pipe \ --name my-pipe \ --endpoint-url $AWS_ENDPOINT_URL # List all pipes aws pipes list-pipes \ --endpoint-url $AWS_ENDPOINT_URL # Start a pipe aws pipes start-pipe \ --name my-pipe \ --endpoint-url $AWS_ENDPOINT_URL # Stop a pipe aws pipes stop-pipe \ --name my-pipe \ --endpoint-url $AWS_ENDPOINT_URL # Update a pipe aws pipes update-pipe \ --name my-pipe \ --target arn:aws:lambda:us-east-1:000000000000:function:new-function \ --endpoint-url $AWS_ENDPOINT_URL # Delete a pipe aws pipes delete-pipe \ --name my-pipe \ --endpoint-url $AWS_ENDPOINT_URL几个值得注意的细节describe-pipe返回完整的管道配置含SourceParameters、TargetParameters、EnrichmentParameters、Tags、CreationTime、LastModifiedTime等见 PipesController.buildPipeResponselist-pipes只返回摘要ARN、名称、源、目标、期望状态、当前状态、时间戳不含SourceParameters与 AWS 行为一致见 buildPipeListEntrylist-pipes支持--name-prefix、--source-prefix、--target-prefix、--desired-state、--current-state过滤对应 PipesService.listPipesupdate-pipe采用“未提供的属性保持原值”的语义只有显式传入的字段才会被覆盖null 属性不会清空既有值见 writePipeConfigurationdelete-pipe会先停止轮询再删除存储记录stop-pipe会取消轮询定时器、清空 Kinesis/DynamoDB 迭代器并关闭 Kafka 消费者见 PipesPoller.stopPolling。Floci 的pipes服务还实现了资源标签能力tagResource/untagResource/listTags并作为pipes:pipe资源类型接入 Resource Explorer见 PipesService.getSupportedResourceTypes。管道状态机DesiredState与CurrentState两个概念需要区分DesiredStateRUNNING或STOPPED表示期望状态CurrentState实际状态。PipeState 枚举完整定义了以下状态STARTING— 管道正在启动RUNNING— 管道正在活跃处理事件STOPPING— 管道正在停止STOPPED— 管道已停止不再处理事件DELETED— 管道已被删除源码中枚举还包含CREATING、UPDATING、DELETING、CREATE_FAILED、UPDATE_FAILED、START_FAILED、STOP_FAILED、DELETE_FAILED等完整状态集与 AWS 线格式保持一致。状态流转逻辑在 PipesService 中创建时未指定DesiredState则默认为RUNNING此时CurrentState直接取RUNNING并启动轮询createPipe创建时指定--desired-state STOPPEDCurrentState即为STOPPED不会启动轮询start-pipe将 DesiredState 与 CurrentState 都置为RUNNING并恢复轮询startPipestop-pipe先停止轮询再将两个状态置为STOPPEDstopPipe。兼容性测试 pipes.bats 中验证了--desired-state STOPPED创建后CurrentState STOPPED的行为。另外Floci 重启后会自动恢复CurrentState RUNNING的管道轮询startPersistedPollers保证模拟环境中断后管道能继续消费。支持的源与目标Floci 模拟的 EventBridge Pipes 支持以下源与目标类型由 PipesPoller.pollAndInvoke 依据源 ARN 中的服务段自动分派Sources源Amazon SQS 队列arn:aws:sqs:...Amazon Kinesis 流arn:aws:kinesis:...Amazon DynamoDB 流arn:aws:dynamodb:...Kafka topics —— MSKarn:aws:kafka:...与自管理smk://协议Targets目标Lambda 函数:lambda:或:function:SQS 队列SNS 主题Kinesis 流Step Functions 状态机:states:目标分派逻辑在 PipesTargetInvoker.invoke按 ARN 特征路由到对应服务Lambda以RequestResponse同步调用若返回FunctionError视为投递失败并转入 DLQinvokeLambdaSQS将 payload 作为消息体发送到队列SNS以SubjectPipes发布到主题invokeSnsEventBridge向事件总线PutEvents默认Sourceaws.pipes、DetailTypePipeForwarded可通过EventBridgeEventBusParameters覆盖invokeEventBridgeStep Functions以随机执行名pipes-uuid启动执行invokeStepFunctions。目标侧还支持TargetParameters.InputTemplate模板中用$.path形式的 JSONPath 占位符投递前替换为从事件中提取的值实现在 applyInputTemplate / extractJsonPath。Kafka 源与 Karapace REST 桥接Kafka 源MSK 或自管理smk://是 Floci Pipes 中最特殊的一条路径它按需启动一个 Karapace REST Proxy 边车容器通过 HTTP 轮询 Kafka而不是在 Floci 进程内嵌入kafka-clients参见 KarapaceManager 的注释关联 issue #2916。工作原理每个不同的bootstrap.servers目标对应一个 Karapace 容器多个读取同一集群的 Pipe共享这一个实例首次使用时ensureStarted按需启动通过 SHA-256 摘要的前 8 个十六进制字符生成目标键并用该键构造容器名pipes-kafka-rest-hash容器以KARAPACE_*环境变量配置KARAPACE_HOST0.0.0.0、KARAPACE_PORT8082、KARAPACE_BOOTSTRAP_URIbootstrap、KARAPACE_ADVERTISED_PROTOCOLhttp并显式指定入口/venv/bin/karapace_rest_proxy镜像默认入口是无操作的python3见 KarapaceManager.startContainer启动后通过GET /_health探测就绪30 秒超时随后 Floci 的 PipesKafkaRestClient 使用 Kafka REST Proxy v2 API 子集消费并提交偏移量地址改写若bootstrap.servers是localhost/127.0.0.1会改写为 Docker 主机可达地址host.docker.internal别名并在 Linux 原生 Docker 上显式注入 host-gateway 映射见 resolveForContainer进程退出时PreDestroy统一停止并删除全部 Karapace 容器。相关配置环境变量默认值说明FLOCI_SERVICES_PIPES_ENABLEDtrue启用或禁用该服务FLOCI_SERVICES_PIPES_KAFKA_REST_BRIDGE_DEFAULT_IMAGEghcr.io/aiven-open/karapace:latestKafka 源管道按需启动的 Karapace REST Proxy 边车镜像FLOCI_SERVICES_PIPES_KAFKA_REST_BRIDGE_HOST_PORT_BASE9500分配给 Karapace 边车的主机端口段起点FLOCI_SERVICES_PIPES_KAFKA_REST_BRIDGE_HOST_PORT_MAX9599分配给 Karapace 边车的主机端口段终点配置实现在 EmulatorConfig.PipesServiceConfig / KafkaRestBridgeConfigKafkaRestBridgeConfig定义了默认镜像与端口段KarapaceManager 通过PortAllocator在[9500, 9599]区间分配主机端口容器内运行时不绑定主机端口仅暴露 8082。重要前提带 Kafka 源的管道依赖 Docker 环境因为需要启动 Karapace 边车其余源SQS/Kinesis/DynamoDB Streams不需要 Docker。Kafka 源的参数要求创建 Kafka 源管道时SourceParameters必须携带对应的参数块否则返回ValidationException见 validateSourceConfigurationsmk://前缀 → 必须提供SelfManagedKafkaParameters:kafka:ARN → 必须提供ManagedStreamingKafkaParameters且参数块内TopicName为必填项。示例aws pipes create-pipe \ --name my-kafka-pipe \ --source arn:aws:kafka:us-east-1:000000000000:cluster/my-msk/b4f3c2a1-0000-4000-8000-000000000000-9 \ --target arn:aws:lambda:us-east-1:000000000000:function:my-function \ --role-arn arn:aws:iam::000000000000:role/pipe-role \ --source-parameters { ManagedStreamingKafkaParameters: { TopicName: orders, StartingPosition: TRIM_HORIZON } } \ --endpoint-url $AWS_ENDPOINT_URL自管理 Kafka 的smk://源同理但使用SelfManagedKafkaParameters并在其中给出BootstrapServers。ParallelizationFactor参数语义与当前限制CreatePipe与UpdatePipe在KinesisStreamParameters和DynamoDBStreamParameters源参数块中接受1 到 10 之间的整数ParallelizationFactor与 AWS 线格式一致。DescribePipe会在SourceParameters中原样回显ListPipes只返回摘要、不含SourceParameters与 AWS 一致。aws pipes create-pipe \ --name kinesis-pipe \ --source arn:aws:kinesis:us-east-1:000000000000:stream/events \ --target arn:aws:lambda:us-east-1:000000000000:function:my-function \ --role-arn arn:aws:iam::000000000000:role/pipe-role \ --source-parameters {KinesisStreamParameters:{StartingPosition:TRIM_HORIZON,ParallelizationFactor:4}} \ --endpoint-url $AWS_ENDPOINT_URL校验逻辑严格对齐 AWSvalidateParallelizationFactor值超出 1~10 →ValidationException字段出现在与源不匹配的参数块上例如 SQS 源携带KinesisStreamParameters.ParallelizationFactor→ValidationException非数字、非整数、超出 int 范围的宽值如2^64 5也会被拒绝——源码通过canConvertToInt()防止 BigInteger 溢出绕过边界检查。Enforcement status重要限制ParallelizationFactor目前仅持久化并在线格式上回显轮询器尚未按分片并发处理批次。Floci 对每个 Kinesis / DynamoDB Streams 管道只在单个分片shardId-000000000000上打开迭代器、每次只投递一个批次无论配置值是多少见 initKinesisIterator 与 initDynamoDbIterator。多分片轮询与真正的分片级并发被列为后续跟进项。Filter事件过滤与匹配语义管道支持 AWS 风格的FilterCriteria过滤位于SourceParameters.FilterCriteria.Filters每个过滤条件包含PatternJSON 匹配模式。实现位于 PipesFilterMatcher多条 filter 之间是OR关系任一匹配即通过单条 filter 内的多个字段是AND关系模式值必须是数组或对象标量模式会被判定为非法并拒绝该记录支持字符串精确匹配、数字精确匹配、null字段不存在、对象嵌套支持 AWS 特殊模式{prefix: ...}、{suffix: ...}、{equals-ignore-case: ...}、{anything-but: ...}含anything-butprefix、{exists: true|false}、{numeric: [, 5, , 1]}等。SQS 路径下未匹配过滤条件的消息会被立即从源队列删除AWS 行为被过滤掉的消息视为已消费见 PipesPoller.pollSqs。Kafka 路径下未匹配的记录会推进提交偏移量相当于消费掉见 deliverKafkaRecords。EnrichmentLambda 富化与 DLQ 兜底管道的可选富化步骤source → filter → enrichment → target目前仅为 Lambda enrichment 提供模拟过滤后的批次被同步调用RequestResponse富化函数的响应成为目标的输入。实现在 PipesTargetInvoker.applyEnrichment。行为细则与 AWS 对齐空响应跳过目标空 body、null、{}、[]都会消费源记录而不调用目标非空数组如[{}]仍会调用目标携带一个空 payload 元素。注意源码会先解析 JSON 以区分{ }、[ ]等空白变体Lambda 富化返回FunctionError会判定批次失败源记录被路由到管道的死信队列DLQ而非静默消费非 Lambda 富化类型API destinations、API Gateway、Step Functions Express在 AWS 上合法但 Floci不模拟配置了此类富化的管道会把批次判失败并转入 DLQ而不是投递未富化的 payload富化目前只在 SQS 源路径生效Kinesis、DynamoDB Streams 与 Kafka 源会把过滤后的记录直接投递给目标配置在这些源上的 enrichment 暂不生效见 PipesPoller.deliverRecords 的注释。富化响应到目标的形状转换富化返回的是单个 JSON 值非数组时代表一条事件会被包装成单元素数组返回数组则视为已是批次、原样转发。该逻辑在 asEventArray。同时Lambda 目标接收的输入保持{Records: [...]}批次形状而 Step Functions、SQS、SNS、EventBridge 目标接收富化响应的原始内容数组包装会破坏它们的输入见 deliverEnrichedBatch。DLQ死信队列配置管道的 DLQ 通过源参数块中的DeadLetterConfig.Arn配置支持SqsQueueParameters、KinesisStreamParameters、DynamoDBStreamParameters三个块getDlqArnaws pipes create-pipe \ --name my-pipe \ --source arn:aws:sqs:us-east-1:000000000000:source-queue \ --target arn:aws:lambda:us-east-1:000000000000:function:my-function \ --role-arn arn:aws:iam::000000000000:role/pipe-role \ --source-parameters { SqsQueueParameters: { BatchSize: 10, DeadLetterConfig: {Arn: arn:aws:sqs:us-east-1:000000000000:dlq-queue} } } \ --endpoint-url $AWS_ENDPOINT_URL投递失败时sendToDeadLetterQueue 会把失败 payload 发送到 DLQ若未配置 DLQ则返回投递失败计数并记录告警日志。批次大小与轮询行为默认BatchSize为 10可通过各源参数块的BatchSize字段调整getBatchSize轮询间隔固定为 1 秒POLL_INTERVAL_MS 1000通过 Vert.x 周期定时器驱动并使用独立线程池执行保证同一管道不会并发轮询activePolls防重入见 PipesPollerKinesis 迭代器过期ExpiredIteratorException或 DynamoDB 流数据被裁剪TrimmedDataAccessException时会自动重置迭代器并重新拉取pollKinesis、pollDynamoDbStreamsKafka 路径下REST Proxy 一次轮询可能返回超过一个BatchSize的记录因此一个轮询周期会顺序投递多个批次一旦某分区投递失败该分区在后续批次中也不会提交偏移量保证“先失败记录重投”的顺序语义deliverKafkaBatch。与 CloudFormation 的集成Pipes 服务实现了ResourceProvider接口支持 CloudFormation 对AWS::Pipes::Pipe资源的创建、更新与删除提供器位于 cloudformation/provisioners 下的PipesCfnProvisioner并配套了回滚恢复测试 PipesCfnRollbackRestoresExactConfigurationTest.java 与 PipesCfnProvisionerTest.java。restorePipe路径采用“null 即清空”的全量恢复语义确保 CloudFormation 回滚时能把管道精确还原到快照状态PipesService.restorePipe。测试与验证仓库提供了多层测试覆盖 pipes 服务单元测试PipesServiceTest.java、PipesPollerTest.java、PipesTargetInvokerTest.java、PipesFilterMatcherTest.java、KarapaceManagerTest.java、PipesKafkaRestClientTest.java集成测试PipesIntegrationTest.java、PipesPollerIntegrationTest.javaAWS CLI 兼容性测试pipes.bats 覆盖 create/describe/list/update/start/stop 全流程SDK 各语言套件sdk-test-go、sdk-test-node、sdk-test-java、sdk-test-python中的pipes用例验证多语言 SDK 兼容性。结语Floci 的 EventBridge Pipes 实现覆盖了从创建、启停、更新、删除到列表过滤的完整生命周期并对 SQS/Kinesis/DynamoDB Streams/Kafka 四种源与 Lambda/SQS/SNS/Kinesis/Step Functions/EventBridge 多种目标提供事件搬运能力Kafka 源通过按需启动的 Karapace REST 边车实现零内嵌依赖的轮询。在使用时请特别注意当前仿真边界ParallelizationFactor仅为线格式兼容单分片顺序轮询、enrichment 仅支持 Lambda 且仅作用于 SQS 源路径、Kafka 源依赖 Docker。这些限制都以代码注释和文档!!! note块的形式明确标注便于你在搭建本地事件管道时做出准确预期。【免费下载链接】flociLight, fluffy, and always free - The AWS Local Emulator alternative项目地址: https://gitcode.com/gh_mirrors/fl/floci创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表