You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

启用Wait for Callback后Step Function卡在Glue Crawler步骤的问题

问题原因分析

你使用的arn:aws:states:::aws-sdk:glue:startCrawler.waitForTaskToken属于回调等待模式,该模式要求必须通过调用SendTaskSuccess/SendTaskFailure API传入任务令牌(Task Token),Step Functions才能继续执行后续步骤。但Glue Crawler本身没有内置机制在运行完成后自动发送这个回调请求,因此即使Crawler执行完毕,Step Functions会一直处于等待状态,无法推进流程。

解决方案

针对这个问题,有两种可行的解决方式,根据你的场景选择即可:

方案一:轮询等待Crawler完成(推荐,无需额外组件)

修改Step Functions状态机,先启动Crawler,再通过轮询检查Crawler状态,直到其运行完成或失败:

{
  "Comment": "Orchestrate Glue Crawler and Job with Polling",
  "StartAt": "StartCrawler",
  "States": {
    "StartCrawler": {
      "Type": "Task",
      "Next": "WaitBeforeCheck",
      "Parameters": {
        "Name": "farmlands-raw-crawler-sam"
      },
      "Resource": "arn:aws:states:::aws-sdk:glue:startCrawler"
    },
    "WaitBeforeCheck": {
      "Type": "Wait",
      "Seconds": 30,
      "Next": "CheckCrawlerStatus"
    },
    "CheckCrawlerStatus": {
      "Type": "Task",
      "Parameters": {
        "Name": "farmlands-raw-crawler-sam"
      },
      "Resource": "arn:aws:states:::aws-sdk:glue:getCrawler",
      "Next": "IsCrawlerComplete?"
    },
    "IsCrawlerComplete?": {
      "Type": "Choice",
      "Choices": [
        {
          "Variable": "$.Crawler.State",
          "StringEquals": "READY",
          "Next": "Glue StartJobRun"
        },
        {
          "Variable": "$.Crawler.State",
          "StringEquals": "FAILED",
          "Next": "CrawlerFailed"
        }
      ],
      "Default": "WaitBeforeCheck"
    },
    "CrawlerFailed": {
      "Type": "Fail",
      "Cause": "Glue Crawler execution failed",
      "Error": "CrawlerFailure"
    },
    "Glue StartJobRun": {
      "Type": "Task",
      "Resource": "arn:aws:states:::glue:startJobRun.sync",
      "Parameters": {
        "JobName": "fetch-from-azure-etl-sam"
      },
      "End": true
    }
  }
}

逻辑说明:

  1. 调用startCrawler直接启动爬虫,无需等待
  2. 设置固定等待时长(可根据爬虫实际运行时长调整)
  3. 调用getCrawler获取当前爬虫状态
  4. 通过Choice分支判断:
    • 状态为READY:表示爬虫已完成,进入后续Job步骤
    • 状态为FAILED:标记流程失败
    • 其他状态:继续等待并轮询

方案二:EventBridge+Lambda实现回调(适合精确触发场景)

如果坚持使用waitForTaskToken模式,需要额外搭建回调机制:

1. 封装Crawler启动逻辑(Lambda)

编写Lambda函数,负责启动Crawler并将Task Token存储到DynamoDB(或其他存储):

import boto3
import json

glue = boto3.client('glue')
dynamodb = boto3.client('dynamodb')

def lambda_handler(event, context):
    crawler_name = event['CrawlerName']
    task_token = event['TaskToken']
    
    # 启动Glue Crawler
    glue.start_crawler(Name=crawler_name)
    
    # 将Task Token与Crawler名称关联存储
    dynamodb.put_item(
        TableName='CrawlerTaskTokens',
        Item={
            'CrawlerName': {'S': crawler_name},
            'TaskToken': {'S': task_token}
        }
    )
    
    return {
        'statusCode': 200,
        'body': json.dumps('Crawler started and task token stored')
    }

2. 配置EventBridge规则

创建EventBridge规则,监听Glue Crawler的CrawlerStateChange事件,当状态变为SUCCEEDED或FAILED时触发回调Lambda:

  • 事件源:aws.glue
  • 事件类型:CrawlerStateChange
  • 目标:回调Lambda函数

3. 回调Lambda函数

编写Lambda函数,根据EventBridge事件中的Crawler状态,调用Step Functions的SendTaskSuccess/SendTaskFailure:

import boto3
import json

stepfunctions = boto3.client('stepfunctions')
dynamodb = boto3.client('dynamodb')

def lambda_handler(event, context):
    crawler_name = event['detail']['crawlerName']
    crawler_state = event['detail']['state']
    
    # 从DynamoDB获取对应Task Token
    response = dynamodb.get_item(
        TableName='CrawlerTaskTokens',
        Key={'CrawlerName': {'S': crawler_name}}
    )
    task_token = response['Item']['TaskToken']['S']
    
    if crawler_state == 'SUCCEEDED':
        stepfunctions.send_task_success(
            taskToken=task_token,
            output=json.dumps({"CrawlerName": crawler_name, "State": crawler_state})
        )
    else:
        stepfunctions.send_task_failure(
            taskToken=task_token,
            error="CrawlerFailed",
            cause=f"Glue Crawler {crawler_name} failed with state {crawler_state}"
        )
    
    return {
        'statusCode': 200,
        'body': json.dumps('Callback sent successfully')
    }

4. 修改Step Functions状态机

{
  "Comment": "Orchestrate Glue Crawler with Callback",
  "StartAt": "StartCrawlerWithToken",
  "States": {
    "StartCrawlerWithToken": {
      "Type": "Task",
      "Resource": "arn:aws:states:::lambda:invoke",
      "Parameters": {
        "FunctionName": "YourStartCrawlerLambdaArn",
        "Payload": {
          "CrawlerName": "farmlands-raw-crawler-sam",
          "TaskToken.$": "$$.Task.Token"
        }
      },
      "Next": "WaitForCallback"
    },
    "WaitForCallback": {
      "Type": "Task",
      "Resource": "arn:aws:states:::aws-sdk:stepfunctions:waitForTaskToken",
      "Next": "Glue StartJobRun"
    },
    "Glue StartJobRun": {
      "Type": "Task",
      "Resource": "arn:aws:states:::glue:startJobRun.sync",
      "Parameters": {
        "JobName": "fetch-from-azure-etl-sam"
      },
      "End": true
    }
  }
}

注意事项:

  • 需为Lambda配置对应权限:Glue启动权限、DynamoDB读写权限、Step Functions回调权限
  • 需提前创建存储Task Token的DynamoDB表

内容的提问来源于stack exchange,提问作者Hilkjh

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 04:50:44