
后端云原生微服务【免费下载链接】kubelessKubernetes Native Serverless Framework项目地址https://gitcode.com/gh_mirrors/ku/kubeless点击查看免费下载导读Kubeless 允许你让任意 Kubeless 函数响应数据流Data Stream中写入的记录从而把流式数据与无服务器函数直接打通。本文以 docs/streaming-functions.md 为骨架完整讲解在 Kubernetes 集群中部署 Kinesis 触发器控制器、通过 Kubernetes Secret 注入 AWS 凭证、用kubeless trigger kinesis创建触发器并向流中发布记录的全过程。读完本文你将掌握 Kinesis 触发器从部署、创建、发布到日志验证的完整实战链路并理解其底层 CRD 与 CLI 实现原理。说明Kubeless 目前支持的流式数据源为AWS Kinesis。本文所有命令、参数与字段均以当前仓库源码为准。一、数据流事件在 Kubeless 中的定位在 Kubeless 中事件源Event Source与运行时Runtime是解耦的触发器Trigger负责把外部事件转成对函数的 HTTP 调用。数据流事件是其中一类当记录Record被写入流时触发器控制器把记录投递到关联函数的 HTTP 端点函数随后执行。其核心链路为Kinesis Stream持续接收上游写入的记录Kinesis Trigger Controller监听该流对应的kinesistriggers.kubeless.ioCRD 资源控制器消费到记录后以 HTTP POST 方式调用关联的 Kubeless 函数函数处理事件并返回结果日志中可看到收到的原始事件负载。这个事件 → HTTP 调用 → 函数执行的模型与 HTTP、CronJob、Kafka、NATS 等触发器保持一致是 Kubeless 触发器体系中的统一范式。二、第一步部署 Kinesis Trigger Controller要触发函数响应 Kinesis 流中的记录首先需要在集群中部署Kubeless AWS Kinesis Trigger Controller。该控制器由配套项目kubeless/kinesis-trigger提供通过其发布的最新 release 清单安装export RELEASE$(curl -s https://api.github.com/repos/kubeless/kinesis-trigger/releases/latest | grep tag_name | cut -d -f 4) kubectl create -f https://github.com/kubeless/kinesis-trigger/releases/download/$RELEASE/kinesis-$RELEASE.yaml这段命令先通过 GitHub API 获取kinesis-trigger项目的最新 release 版本号再下载并应用对应的kinesis-$RELEASE.yaml清单。该清单会同时创建Kinesis 触发器控制器 Deployment和kinesistriggers.kubeless.ioCRD。部署完成后在kubeless命名空间下可以看到控制器 Pod 正常运行$ kubectl get pods -n kubeless NAME READY STATUS RESTARTS AGE kinesis-trigger-controller-65c78f9f44-v5flq 1/1 Running 0 1h kubeless-controller-manager-6b7cdcdc76-x6gsd 1/1 Running 0 13h同时集群中会出现新的 CRD 资源类型kinesistriggers.kubeless.io$ kubectl get crd NAME AGE cronjobtriggers.kubeless.io 13h functions.kubeless.io 13h httptriggers.kubeless.io 13h kinesistriggers.kubeless.io 13hkinesistriggers.kubeless.io与functions.kubeless.io、httptriggers.kubeless.io、cronjobtriggers.kubeless.io并列存在说明 Kinesis 触发器是 Kubeless 触发器体系中一等公民CRD 化的资源类型。从当前仓库源码看Kubeless CLI 的 Kinesis 命令正是通过引用kinesis-trigger项目的 clientset 来读写这一 CRD 资源见 cmd/kubeless/trigger/kinesis/kinesis_trigger.go 中github.com/kubeless/kinesis-trigger/pkg/...的导入这也是为什么需要先安装上述清单的原因。三、认识kubeless trigger kinesis命令家族Kubeless CLI 内置了对 Kinesis 类型触发器的完整生命周期管理。执行帮助命令可以看到所有可用子命令$ kubeless trigger kinesis --help kinesis trigger command allows users to create, list, update, delete Kinesis triggers running on Kubeless Usage: kubeless trigger kinesis SUBCOMMAND [flags] kubeless trigger kinesis [command] Available Commands: create Create a Kinesis trigger create-stream Create a Kinesis stream delete Delete a Kinesis trigger list list all Kinesis triggers deployed to Kubeless publish publish message to a Kinesis stream update Update a Kinesis trigger Flags: -h, --help help for kinesis Use kubeless trigger kinesis [command] --help for more information about a command.六个子命令覆盖了从建流到发布再到触发器增删改查的完整闭环子命令作用对应源码create创建 Kinesis 触发器关联函数与流create.gocreate-stream创建 Kinesis 流为快速测试提供的便捷命令stream_create.golist列出已部署到 Kubeless 的所有 Kinesis 触发器list.goupdate更新触发器支持局部更新update.godelete删除触发器delete.gopublish向 Kinesis 流发布消息publish.go从 kinesis_trigger.go 的init()可以看出这六个子命令在程序启动时统一注册到KinesisTriggerCmd下Run入口直接输出帮助信息具体逻辑全部下沉到各子命令实现。四、准备 AWS 凭证Kubernetes Secret创建 Kinesis 触发器前必须先让 Kubeless 拿到访问 AWS Kinesis 流的凭证。Kubeless 采用Kubernetes Secret在集群内保存凭证控制器运行时会读取该 Secret 来访问流。Secret 需要包含两个键aws_access_key_idAWS 访问密钥 IDaws_secret_access_keyAWS 访问密钥。通常你的密钥存在于~/.aws/credentials使用 AWS CLI 时也可以从 AWS 控制台创建 Access Keys。创建 Secret 的命令如下kubectl create secret generic ec2 --from-literalaws_access_key_id$AWS_ACCESS_KEY_ID --from-literalaws_secret_access_key$AWS_SECRET_ACCESS_KEY这里的 Secret 名为ec2可自定义之后创建触发器时通过--secret引用。注意两点键名必须是aws_access_key_id和aws_secret_access_key。从 publish.go 的源码可以看到CLI 在发布消息前会显式校验secret.Data中是否存在这两个键缺失会直接Fatalf终止Secret 与触发器应在同一命名空间。create命令通过cli.Core().Secrets(ns).Get(secretName, ...)验证 Secret 存在ns默认取 Kubeless 默认命名空间也可用-n指定。五、创建 Kinesis 触发器参数全解析凭证就绪后即可创建触发器并把一个 Kubeless 函数与 Kinesis 流关联起来kubeless trigger kinesis create test-trigger --function-name post-python --aws-region us-west-2 --shard-id shardId-000000000000 --stream my-kinesis-stream --secret ec2各参数含义如下与 create.go 中注册的 Flag 一一对应参数必填说明trigger_name是触发器名称位置参数如test-trigger--function-name是要关联的 Kubeless 函数名创建前会校验该 Function CRD 是否存在--aws-region是Kinesis 流所在的 AWS 区域如us-west-2--shard-id是记录所放置的分片 ID如shardId-000000000000--stream是Kinesis 流的名称--secret是存放 AWS 凭证的 Kubernetes Secret 名-n, --namespace否触发器所在命名空间默认取 Kubeless 默认命名空间--endpoint否覆盖 AWS 默认服务 URL可用于对接本地 Kinesis 模拟服务--dryrun否只输出触发器 YAML/JSON 清单而不实际创建-o, --output否与--dryrun配合的清单输出格式默认yaml--shard-id的获取方式分片 ID 需要从流的描述信息中取得。Kinesis 流被拆分为多个分片Shard记录根据分区键Partition Key的哈希被分配到特定分片因此创建触发器时要指明监听哪个分片$ aws kinesis describe-stream --stream-name my-kinesis-stream { StreamDescription: { RetentionPeriodHours: 24, StreamName: my-kinesis-stream, Shards: [ { ShardId: shardId-000000000000, HashKeyRange: { EndingHashKey: 340282366920938463463374607431768211455, StartingHashKey: 0 }, SequenceNumberRange: { StartingSequenceNumber: 49584495912138607235774073050889122383423872293029281794 } } ], StreamARN: arn:aws:kinesis:us-west-2:159706291352:stream/my-kinesis-stream, EnhancedMonitoring: [ { ShardLevelMetrics: [] } ], StreamStatus: ACTIVE } }重点关注Shards[].ShardId即--shard-id参数值与StreamStatus状态须为ACTIVE才能正常写入记录。关于参数校验从 create.go 可看到--stream、--aws-region、--shard-id、--function-name、--secret五个参数均通过MarkFlagRequired标记为必填同时创建前会先调用GetFunctionCustomResource校验目标函数已存在、调用 Kubernetes API 校验 Secret 已存在任一步失败都会终止创建。六、验证触发器查看 CRD 对象创建成功后可以像查看普通 Kubernetes 资源一样查看触发器对象$ kubectl get kinesistriggers.kubeless.io test -o yaml apiVersion: kubeless.io/v1beta1 kind: KinesisTrigger metadata: labels: created-by: kubeless name: test namespace: default spec: aws-region: us-west-2 function-name: post-python secret: ec2 shard: shardId-000000000000 stream: my-kinesis-stream这个 YAML 与create命令的入参一一对应spec.function-name是关联函数spec.stream是流名spec.aws-region是区域spec.shard是分片 IDspec.secret是凭证 Secret。metadata.labels中created-by: kubeless标记由 CLI 创建见 create.go 中对ObjectMeta.Labels的设置。提示如果不想直接落地资源可先运行kubeless trigger kinesis create ... --dryrun -o yaml把将要创建的清单打印出来检查确认无误后再去掉--dryrun正式创建。七、向流发布记录两种方式触发器就绪后即可向流中写入记录来触发函数。Kubeless CLI 提供了publish子命令kubeless trigger kinesis publish --aws-region us-west-2 --secret ec2 --partition-key 123 --stream my-kinesis-stream --records hello world从当前仓库源码 publish.go 看publish子命令实际使用的参数为参数必填说明--stream是目标 Kinesis 流名称--aws-region是流所在区域--partition-key是分区键决定记录写入哪个分片--records是要发布的记录列表可重复传入多条StringArray类型--secret是AWS 凭证 Secret 名-n, --namespace否Secret 所在命名空间--endpoint否覆盖 AWS 默认服务 URL本地模拟场景实现上publish会从 Secret 中取出aws_access_key_id与aws_secret_access_key通过 AWS SDK 构建带静态凭证的 session再调用PutRecords批量写入——每条记录都携带同一个partition-key见 publish.go 中PutRecordsRequestEntry的组装逻辑。也可以直接使用 AWS CLI 向流写入记录aws kinesis put-record --stream-name my-kinesis-stream --partition-key 123 --data testdata1 aws kinesis put-record --stream-name my-kinesis-stream --partition-key 123 --data testdata2 aws kinesis put-record --stream-name my-kinesis-stream --partition-key 123 --data testdata3连续写入多条记录可以让触发器控制器按序消费便于观察函数被触发的次数与内容。八、验证函数收到的事件Pod 日志写入记录后查看与触发器关联的函数 Pod 日志可以看到每条被消费的记录都以 HTTP POST 形式到达函数$ kubectl logs post-python-59f7fc4b54-4nhbb Bottle v0.12.13 server starting up (using CherryPyServer())... Listening on http://0.0.0.0:8080/ Hit Ctrl-C to quit. {event-time: 2018-05-18 05:40:42.881137473 0000 UTC, extensions: {request: LocalRequest: POST http://post-python.default.svc.cluster.local:8080/}, event-type: application/x-www-form-urlencoded, event-namespace: kinesistriggers.kubeless.io, data: testdata12, event-id: bDRMSN3NPC81ktU} 172.17.0.7 - - [18/May/2018:05:40:42 0000] POST / HTTP/1.1 200 10 Go-http-client/1.1 0/11758 {event-time: 2018-05-18 05:40:44.891994208 0000 UTC, extensions: {request: LocalRequest: POST http://post-python.default.svc.cluster.local:8080/}, event-type: application/x-www-form-urlencoded, event-namespace: kinesistriggers.kubeless.io, data: testdata22, event-id: uHdiWN-lzeKYQyQ} 172.17.0.7 - - [18/May/2018:05:40:44 0000] POST / HTTP/1.1 200 10 Go-http-client/1.1 0/8983 {event-time: 2018-05-18 05:40:45.878361324 0000 UTC, extensions: {request: LocalRequest: POST http://post-python.default.svc.cluster.local:8080/}, event-type: application/x-www-form-urlencoded, event-namespace: kinesistriggers.kubeless.io, data: testdata32, event-id: sRRjSasGVApy8tA}结合日志与 HTTP 请求行POST / HTTP/1.1可以提炼出 Kinesis 触发器的事件负载Event Payload结构字段示例说明event-time2018-05-18 05:40:42.881137473 0000 UTC事件发生时间RFC3339 纳秒精度event-typeapplication/x-www-form-urlencoded事件内容类型控制器以表单编码把记录投递给函数event-namespacekinesistriggers.kubeless.io事件来源命名空间标识事件出自 Kinesis 触发器event-idbDRMSN3NPC81ktU事件唯一 ID便于追踪单条记录datatestdata12流中记录的实际负载数据extensions.requestPOST http://post-python.default.svc.cluster.local:8080/函数被调用时的 HTTP 请求对象可见目标函数是post-python的 ClusterIP Service每条记录都会触发一次独立的函数调用每条日志后紧跟一行POST / HTTP/1.1 200的访问记录说明触发器控制器按记录粒度把流数据转成了函数调用这就是流式函数的核心语义。event-namespace为kinesistriggers.kubeless.io也印证了事件与 CRD 类型的对应关系。九、触发器的生命周期管理除创建外kubeless trigger kinesis还提供完整的运维能力1. 列出触发器kubeless trigger kinesis list输出为表格形式。从 list.go 源码可见列头为NAME、NAMESPACE、REGION、STREAM、SHARD、FUNCTION NAME即每个触发器的关键属性一览NAME NAMESPACE REGION STREAM SHARD FUNCTION NAME test-trigger default us-west-2 my-kinesis-stream shardId-000000000000 post-python2. 更新触发器当流、区域、分片或关联函数发生变化时使用update进行局部更新源码 update.go 中仅当对应 Flag 非空时才覆盖Spec字段未指定的字段保持不变kubeless trigger kinesis update test-trigger --function-name new-function3. 删除触发器kubeless trigger kinesis delete test-trigger删除后触发器对象随之移除函数将不再响应该流的记录。4. 便捷建流create-streamcreate-stream是为快速测试提供的便捷命令直接从 CLI 创建 Kinesis 流stream_create.gokubeless trigger kinesis create-stream my-stream --aws-region us-west-2 --shard-count 1 --secret ec2--shard-count默认值为1即创建单分片流。完整使用示例可参考仓库集成测试 tests/integration-tests-kinesis.bats其中覆盖了创建流、创建触发器、发布记录、验证函数执行等端到端流程。十、源码级原理CLI 与 kinesis-trigger 项目的协作从代码结构上可以清晰看到 Kinesis 触发器在 Kubeless 中的实现边界CLI 层本仓库cmd/kubeless/trigger/kinesis/下的命令实现负责参数解析、校验与资源编排但不实现控制器逻辑CRD 类型与客户端kinesis-trigger 项目CLI 通过导入github.com/kubeless/kinesis-trigger/pkg/apis/kubeless/v1beta1和github.com/kubeless/kinesis-trigger/pkg/utils来构造、读写KinesisTrigger资源控制器kinesis-trigger 项目部署清单安装的kinesis-trigger-controller负责真正消费流并调用函数。因此Kinesis 触发器采用的是核心仓库提供 CLI 独立项目提供控制器/CRD的扩展模式。这也解释了为什么必须先部署kinesis-$RELEASE.yaml才能使用kubeless trigger kinesis create——CLI 在创建前需要 CRD 已注册、控制器已运行。从推理角度可以认为只要存在实现了同样 CRD 的控制器该 CLI 的调用链即可复用。另外create与update命令都支持--dryrun与--output其实现调用kubelessUtils.DryRunFmt将KinesisTrigger对象序列化为 YAML/JSON 打印而不会真正写集群——这在 CI 或脚本化审计场景中非常有用。十一、本地联调用 kinesalite 模拟 Kinesis仓库的 manifests/kinesis/kinesalite.yaml 提供了一份基于kinesaliteKinesis 的本地模拟实现的清单可用于没有真实 AWS 账号时的本地联调kubectl apply -f manifests/kinesis/kinesalite.yaml该清单包含两部分一个NodePort类型的 Servicekinesis暴露端口4567一个 Deploymentkinesis镜像saikocat/kinesalite:1.11.5以--port4567启动。结合 CLI 中普遍存在的--endpoint参数Override AWSs default service URL with the given URL本地联调时可让create、publish、create-stream都指向 kinesalite 服务地址从而在不触碰真实 AWS 的情况下走通建流 → 发布 → 触发函数的完整链路。这正是--endpoint参数存在的核心使用场景。十二、注意事项与适用前提凭证安全AWS 密钥以 Secret 明文存储于集群请按 Kubernetes 最小权限原则管理 Secret 的访问控制避免把密钥写入镜像或提交到代码库分片与分区键触发器监听的是指定--shard-id分片发布记录时必须使用能路由到该分片的partition-key否则记录会进入其他分片而无法被该触发器消费流状态describe-stream返回的StreamStatus必须为ACTIVE才能写入记录刚创建的流需要等待其变为 ACTIVE命名空间一致性Secret、函数、触发器三者建议处于同一命名空间避免跨命名空间引用带来的权限与解析问题必填参数校验create对stream、aws-region、shard-id、function-name、secret五个参数强制必填且会预校验函数与 Secret 是否存在任何一项不满足都会直接报错退出版本演进原文档中的publish示例使用--message参数而当前仓库源码实际使用--records可重复的 StringArray请以仓库实现为准若你的 CLI 版本较早可先执行kubeless trigger kinesis publish --help确认当前参数名。通过本文的完整流程你已经可以部署 Kinesis 触发器控制器 → 用 Secret 注入 AWS 凭证 → 创建并管理触发器 → 向流发布记录 → 在函数日志中验证事件到达。这套数据流 → 触发器 → 函数的链路是 Kubeless 处理流式数据、构建事件驱动应用的入口能力可进一步与 Kafka、NATS 等触发器对比选择最适合业务的事件源。赞分享后端云原生微服务【免费下载链接】kubelessKubernetes Native Serverless Framework项目地址https://gitcode.com/gh_mirrors/ku/kubeless点击查看免费下载相关推荐Kubeless自定义触发器开发从零实现Kinesis触发器Kubeless自定义触发器开发从零实现Kinesis触发器 Kubeless是一个强大的Kubernetes原生无服务器框架允许开发者在Kubernete后端云原生微服务在 Floci 本地模拟器中使用 AWS Kinesis流式数据管道实战指南在 Floci 本地模拟器中使用 AWS Kinesis流式数据管道实战指南 Floci 是一款轻量、免费的开源 AWS 本地模拟器AWS Local EmAWS SDK for PHP v3 操作 Amazon Kinesis 数据流实战指南php/example_code/kinesisAWS SDK for PHP v3 操作 Amazon Kinesis 数据流实战指南php/example_code/kinesis Amazon Ki示例工程教程后端上一篇PreciseRoIPooling社区资源汇总教程、案例、论文与最佳实践下一篇Servest Cookie与Session管理安全用户认证的最佳实践创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考