terraform-provider-aws 的 Amazon MSK Connect 测试数据解析:test-fixtures 目录与 Simulator Connector 实战
【免费下载链接】terraform-provider-awsThe AWS Provider enables Terraform to manage AWS resources.项目地址: https://gitcode.com/GitHub_Trending/te/terraform-provider-aws
导读
internal/service/kafkaconnect/test-fixtures/目录是 AWS Provider 中 Amazon MSK Connect(MSK Connect)资源与数据源验收测试(Acceptance Tests)的测试数据仓库,核心资产是一份检入版本库的 Kafka Connect Simulator 插件 ZIP 包。本文基于 test-fixtures/README.md 展开,结合custom_plugin_test.go、connector_test.go等测试源码与实现,讲清楚这份测试数据在 MSK Connect 插件上传、自定义插件注册、连接器(Connector)创建这条完整链路中的真实用途,并给出可复用的配置要点与测试运行方法。
一、目录定位:为验收测试提供可复现的插件载体
在 AWS Provider 的internal/service/kafkaconnect服务包中,test-fixtures/README.md明确指出:
本目录包含 Amazon MSK Connect 资源与数据源验收测试所用的测试数据;检入的 ZIP 文件包含 Simulator Connector 的 JAR 文件。
目录当前只有两个文件:
README.md:目录说明文档;jcustenborder-kafka-connect-simulator-0.1.120.zip:约 8.6 MB 的插件分发包。
把真实插件二进制"检入(checked-in)"仓库而不是在测试时从网络下载,保证了验收测试的确定性:CI 与本地测试不依赖外网可达性、版本漂移或第三方分发渠道可用性,测试数据与测试代码天然保持版本一致。这种"测试夹具随源码入库"的做法也是该仓库在多服务验收测试中普遍采用的策略。
二、为什么选 Simulator Connector:测试数据生成器
Simulator Connector(kafka-connect-simulator)是一个面向 Kafka Connect 生态的测试数据生成插件,可从 ZIP 内manifest.json确认其元信息:
- 名称:
kafka-connect-simulator,版本0.1.120,发布于2019-08-15; - 描述:A Kafka Connect connector for generating test data(用于生成测试数据的连接器);
- 组件类型:
sink与source两种 connector 均支持; - 许可证:Apache License 2.0。
对验收测试而言,它有两个关键优势:一是纯软件模拟,不依赖任何外部系统(如数据库、SaaS),适合在隔离的 MSK 集群上稳定运行;二是同时具备 sink 与 source 两种形态,覆盖了连接器测试的两类典型配置。因此测试代码可以放心地把它作为aws_mskconnect_custom_plugin与aws_mskconnect_connector两个资源的默认插件载体。
三、ZIP 内部结构:一次真实的解包观察
对jcustenborder-kafka-connect-simulator-0.1.120.zip执行解包可以确认,其内部遵循 Confluent Hub 式插件的标准布局:
jcustenborder-kafka-connect-simulator-0.1.120/ ├── doc/ │ ├── LICENSE │ └── README.md ├── etc/ │ ├── SimulatorSinkConnector.properties # sink 示例配置 │ ├── SimulatorSourceConnector.properties # source 示例配置 │ └── connect-avro-docker.properties ├── lib/ │ ├── kafka-connect-simulator-0.1.120.jar # 主连接器 JAR │ ├── connect-utils-0.4.156.jar │ ├── guava-27.1-jre.jar │ ├── jackson-*(多个) │ ├── jfairy-0.5.3.jar # 假数据生成依赖 │ └── ...(其余传递依赖) └── manifest.json # 插件元数据其中etc/SimulatorSinkConnector.properties给出了 sink 形态的关键配置骨架:
name=MySinkConnector tasks.max=1 topic=simulator connector.class=com.github.jcustenborder.kafka.connect.simulator.SimulatorSinkConnector值得注意的是,仓库测试并没有照搬这份示例 properties,而是在 Terraform 配置里直接以connector_configuration映射给出自己的值(见下文第四节),这说明该 ZIP 更核心的角色是可上传、可被 MSK Connect 识别的插件归档,具体运行参数由各测试自行注入。
四、测试数据如何被消费:从 test-fixtures 到 S3 到 MSK Connect
这份 ZIP 在验收测试中的流转路径,完整对应了真实用户从 S3 上传插件到创建连接器的操作流程。
4.1 上传到 S3:aws_s3_object以本地文件为 source
在 custom_plugin_test.go 的testAccCustomPluginBaseConfig中,测试先用一个 S3 桶 + 版本控制配置(可选)+ S3 对象搭建基础环境:
resource "aws_s3_bucket" "test" { bucket = "terraform-acc-test-..." # 实际由 rName 注入 force_destroy = true } resource "aws_s3_bucket_versioning" "test" { bucket = aws_s3_bucket.test.bucket versioning_configuration { status = false ? "Suspended" : "Enabled" # 由测试参数决定 } } resource "aws_s3_object" "test" { bucket = aws_s3_bucket_versioning.test.bucket key = "jcustenborder-kafka-connect-simulator-0.1.120.zip" source = "test-fixtures/jcustenborder-kafka-connect-simulator-0.1.120.zip" }这里source指向的正是test-fixtures目录下的 ZIP 文件——Terraform 在应用配置时会把该文件原样上传到 S3,模拟"用户把插件包放入 S3 桶"这一真实前置动作。
4.2 注册自定义插件:aws_mskconnect_custom_plugin
随后测试创建 MSK Connect 自定义插件(custom_plugin_test.go):
resource "aws_mskconnect_custom_plugin" "test" { name = "terraform-acc-test-..." # 实际由 rName 注入 content_type = "ZIP" location { s3 { bucket_arn = aws_s3_bucket.test.arn file_key = aws_s3_object.test.key } } }从实现源码 custom_plugin.go 可以看到,Create流程会调用 MSK Connect 的CreateCustomPluginAPI,把content_type(ZIP)、location(S3 的bucket_arn/file_key/object_version)与名称组装进请求,随后通过状态机等待插件从CREATING变为ACTIVE,默认创建超时为 10 分钟;若状态进入CREATE_FAILED,还会把 AWS 返回的错误码与消息透出。若启用 S3 桶版本控制并传入object_version,则可精确锁定插件包的具体版本(对应TestAccKafkaConnectCustomPlugin_objectVersion用例)。
4.3 组装连接器:aws_mskconnect_connector引用自定义插件
插件注册完成后,连接器测试(connector_test.go)把自定义插件与 MSK 集群、IAM、VPC 等环境组合起来:
resource "aws_mskconnect_connector" "test" { name = "terraform-acc-test-..." # 实际由 rName 注入 kafkaconnect_version = "2.7.1" capacity { autoscaling { min_worker_count = 1 max_worker_count = 2 } } connector_configuration = { "connector.class" = "com.github.jcustenborder.kafka.connect.simulator.SimulatorSinkConnector" "tasks.max" = "1" "topics" = "t1" } kafka_cluster { apache_kafka_cluster { bootstrap_servers = aws_msk_cluster.test.bootstrap_brokers_tls vpc { security_groups = [aws_security_group.test.id] subnets = aws_subnet.test[*].id } } } plugin { custom_plugin { arn = aws_mskconnect_custom_plugin.test.arn revision = aws_mskconnect_custom_plugin.test.latest_revision } } service_execution_role_arn = aws_iam_role.test.arn depends_on = [aws_iam_role_policy.test, aws_vpc_endpoint.test] }注意connector_configuration中connector.class填写的正是 Simulator 的SimulatorSinkConnector全限定类名——这与 ZIP 内示例 properties 中的类名一致,而topics等参数则由测试自行定义(与 ZIP 内示例的topic=simulator不同)。plugin块通过arn + revision精确引用上一步注册的自定义插件,其中revision取自latest_revision属性,形成"test-fixtures → S3 → custom plugin → connector"的完整闭环。
五、从测试代码反推参数边界与设计要点
结合 connector.go 的 schema 定义,可以提炼出配置参数的边界(这些约束同样适用于生产环境的 Terraform 配置):
| 参数 | 约束/默认值 | 说明 |
|---|---|---|
capacity.0.autoscaling.min/max_worker_count | 整数,范围 1–10 | 与provisioned_capacity通过ExactlyOneOf互斥,二者只能选一 |
capacity.0.autoscaling.mcu_count | 默认 1,取值 {1, 2, 4, 8} | 每个 worker 的计算单元数 |
scale_in_policy.cpu_utilization_percentage | 1–100 | 缩容阈值,如测试中的 25 |
scale_out_policy.cpu_utilization_percentage | 1–100 | 扩容阈值,如测试中的 75 |
capacity.0.provisioned_capacity.worker_count | 整数,范围 1–10 | 固定容量模式必填 |
connector_configuration | TypeMap,键值均为字符串 | 直接透传给 MSK Connect 的 connector 配置 |
kafkaconnect_version | 字符串(测试用2.7.1) | MSK Connect 运行时版本 |
连接器的默认超时在源码中定义为:创建/更新 20 分钟、删除 10 分钟(connector.go)。此外,测试还覆盖了disappears(资源被外部删除后重规划)、tags(标签增删改)、update(从 autoscaling 切换到 provisioned_capacity、修改connector_configuration.topics等)等场景,connector_test.go中的每个resource.TestCheckResourceAttr断言都直接锚定上述测试配置,可作为理解各属性语义的活文档。
六、资源族全景:测试数据之外还有哪些配套资源
test-fixtures服务于internal/service/kafkaconnect服务包下的一组资源与数据源:
aws_mskconnect_custom_plugin(资源 + 数据源):插件注册,测试数据 ZIP 的直接消费者;aws_mskconnect_connector(资源 + 数据源):连接器生命周期管理,引用插件与 worker 配置;aws_mskconnect_worker_configuration(资源 + 数据源):worker 配置(properties 文件内容),与连接器通过worker_configuration块关联(见connector_test.go的testAccConnectorConfig_allAttributes)。
其中custom_plugin_data_source.go与connector_data_source.go分别提供按名称/ARN 查询的读取能力,sweep.go提供了清理残留测试资源的 sweeper。整个服务包还遵循仓库统一的分区与端点校验(PreCheckPartitionHasService、names.KafkaConnectServiceID)模式。
七、如何运行与扩展这些测试
验收测试需要真实的 AWS 凭证与 MSK 相关权限,运行前请先阅读仓库的 运行与编写验收测试指南 与 验收测试环境变量说明,配置AWS_*凭证后执行:
# 运行 MSK Connect 服务包的全部验收测试 make testacc TESTS=TestAccKafkaConnect ./internal/service/kafkaconnect/...或按用例名定向执行,例如:
make testacc TESTS=TestAccKafkaConnectCustomPlugin_basic ./internal/service/kafkaconnect/... make testacc TESTS=TestAccKafkaConnectConnector_basic ./internal/service/kafkaconnect/...若需为其他连接器(如 JDBC、S3 sink)编写新测试,可按同样模式:将插件 ZIP 放入对应服务的test-fixtures/目录,用aws_s3_object上传,再经aws_mskconnect_custom_plugin注册即可,无需改动测试框架本身。
八、小结
internal/service/kafkaconnect/test-fixtures/README.md虽短,却是一个典型"小而关键"的测试基础设施:它以一份检入版本库的 Simulator Connector ZIP 包,支撑起 MSK Connect 三个资源/数据源全部验收测试的插件供给,使测试具备确定性、可复现性与离线可运行性。从本文可以看出,阅读这类 test-fixtures 说明时,最有价值的信息往往不在 README 本身,而在测试代码如何消费这些数据——custom_plugin_test.go与connector_test.go中的配置片段就是最好的实战范例。
【免费下载链接】terraform-provider-awsThe AWS Provider enables Terraform to manage AWS resources.项目地址: https://gitcode.com/GitHub_Trending/te/terraform-provider-aws
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考