如何让Serverless部署的Kinesis Firehose等待Elasticsearch域创建完成?
当然有办法解决这个问题!在Serverless Framework的配置里,我们可以通过资源依赖和自定义等待逻辑来确保Firehose只在ES域完全创建好之后才开始部署。下面是具体的实现方式:
方法一:利用CloudFormation的DependsOn属性(基础方案)
Serverless Framework本质是基于CloudFormation的,所以可以直接在Firehose的资源定义里添加DependsOn属性,明确依赖你的Elasticsearch域资源。这样CloudFormation会自动等待ES域创建完成后再创建Firehose。
比如在你的serverless.yml里:
resources: Resources: MyElasticsearchDomain: Type: AWS::Elasticsearch::Domain Properties: # 你的ES域配置 DomainName: my-elastic-search # ...其他配置 MyFirehoseStream: Type: AWS::KinesisFirehose::DeliveryStream DependsOn: MyElasticsearchDomain # 关键:添加这个依赖 Properties: DeliveryStreamName: MyFirehoseStream ElasticsearchDestinationConfiguration: DomainARN: !GetAtt MyElasticsearchDomain.DomainArn # ...其他Firehose配置
这个方法是最直接的,CloudFormation会自动处理等待逻辑,但要注意:如果你的ES域是通过其他方式创建的(不是在同一个Serverless栈里),那这个方法就不适用了。
方法二:自定义CloudFormation等待条件(跨栈或外部创建的ES域)
如果你的ES域是在另一个CloudFormation栈或者手动创建的,你可以用CloudFormation的AWS::CloudFormation::WaitCondition和AWS::CloudFormation::WaitConditionHandle来手动等待ES域进入可用状态。
步骤大概是:
- 创建一个WaitConditionHandle资源。
- 创建一个Lambda函数,用来检查ES域的状态,当状态变为
ACTIVE时,向WaitConditionHandle发送信号。 - 在Firehose资源里添加
DependsOn到WaitCondition,确保Firehose在等待条件满足后才创建。
示例配置片段:
resources: Resources: WaitForESHandle: Type: AWS::CloudFormation::WaitConditionHandle WaitForESCondition: Type: AWS::CloudFormation::WaitCondition DependsOn: CheckESStatusLambda Properties: Handle: !Ref WaitForESHandle Timeout: '3600' # 最多等待1小时,根据你的ES创建时间调整 CheckESStatusLambda: Type: AWS::Lambda::Function Properties: Runtime: python3.9 Handler: index.lambda_handler Role: !GetAtt CheckESStatusLambdaRole.Arn Environment: Variables: ES_DOMAIN_ARN: arn:aws:es:us-east-1:1234567890:domain/my-elastic-search WAIT_HANDLE: !Ref WaitForESHandle Code: ZipFile: | import boto3 import os import json es_client = boto3.client('es') wait_handle = os.environ['WAIT_HANDLE'] def lambda_handler(event, context): domain_arn = os.environ['ES_DOMAIN_ARN'] domain_name = domain_arn.split('/')[-1] response = es_client.describe_elasticsearch_domain(DomainName=domain_name) status = response['DomainStatus']['DomainStatus'] if status == 'ACTIVE': # 发送成功信号 import requests requests.put(wait_handle, data=json.dumps({'Status': 'SUCCESS'})) else: # 抛出异常让Lambda重试,CloudFormation会等待 raise Exception(f"ES domain status is {status}, not ACTIVE yet") CheckESStatusLambdaRole: Type: AWS::IAM::Role Properties: AssumeRolePolicyDocument: Version: '2012-10-17' Statement: - Effect: Allow Principal: Service: lambda.amazonaws.com Action: sts:AssumeRole Policies: - PolicyName: CheckESStatus PolicyDocument: Version: '2012-10-17' Statement: - Effect: Allow Action: es:DescribeElasticsearchDomain Resource: arn:aws:es:us-east-1:1234567890:domain/my-elastic-search - Effect: Allow Action: logs:CreateLogGroup Resource: !Sub arn:aws:logs:${AWS::Region}:${AWS::AccountId}:log-group:/aws/lambda/CheckESStatusLambda:* - Effect: Allow Action: - logs:CreateLogStream - logs:PutLogEvents Resource: !Sub arn:aws:logs:${AWS::Region}:${AWS::AccountId}:log-group:/aws/lambda/CheckESStatusLambda:*:* MyFirehoseStream: Type: AWS::KinesisFirehose::DeliveryStream DependsOn: WaitForESCondition Properties: # 你的Firehose配置
这个方法更灵活,适合ES域不在当前Serverless栈的情况,但需要额外的Lambda资源来检查状态。
方法三:使用Serverless插件(简化方案)
还有一些Serverless插件可以帮你处理资源等待的问题,比如serverless-plugin-wait-for-stack(如果ES在另一个栈)或者自定义插件,但要注意插件的维护状态。不过如果是同一个栈的话,方法一就足够了,不需要额外插件。
总结一下:
- 如果ES域和Firehose在同一个Serverless栈里,直接用
DependsOn是最简单的方案。 - 如果ES域在外部,就用方法二的等待条件+Lambda检查状态。
内容的提问来源于stack exchange,提问作者twiz

