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

部署Azure Durable Function时提示未找到HTTP触发器

问题:Azure Durable Function部署后无HTTP触发器检测到

我开发了一个Azure Durable Function,本地运行时可以正常从API提取数据并加载到PostgreSQL数据库,但部署到App Service Plan托管的Function App后,出现错误提示“No HTTP triggers found”。

基础配置

  • 托管方式:App Service Plan
  • 触发类型:HTTP触发

核心函数代码

import azure.functions as func
import azure.durable_functions as df
from fluxx_extraction import data_extractor

myApp = df.DFApp(http_auth_level=func.AuthLevel.ANONYMOUS)

# HTTP触发的Durable客户端函数
@myApp.route(route="orchestrators/{functionName}")
@myApp.durable_client_input(client_name="client")
async def Http_trigger(req: func.HttpRequest, client):
    function_name = req.route_params.get('functionName')
    instance_id = await client.start_new(function_name)
    response = client.create_check_status_response(req, instance_id)
    return response

# 编排器函数
@myApp.orchestration_trigger(context_name="context")
def hello_orchestrator(context):
    result1 = yield context.call_activity("activity_function")

    return result1

# 活动函数
@myApp.activity_trigger(input_name="name")
def activity_function(city: str):
    data_extractor()
    return "extraction completed"

data_extractor.py代码

import asyncio
import aiohttp
import asyncpg
import requests
import pandas as pd
from datetime import datetime, timezone
from sqlalchemy import create_engine, text
import logging
from tenacity import retry, stop_after_attempt, wait_fixed, RetryError

# 日志配置
logging.basicConfig(level=logging.INFO)

# API限流配置
REQUESTS_PER_MINUTE = 400
REQUESTS_INTERVAL = 60 / REQUESTS_PER_MINUTE  # 请求间隔(秒)
RECORDS_PER_PAGE = 100  # 每页记录数

# 控制并发请求的信号量
semaphore = asyncio.Semaphore(REQUESTS_PER_MINUTE // 60)  # 每秒并发请求数

async def fetch_access_token(api_credentials):
    token_url = "https://appi.fluxx.io/oauth/token"
    try:
        response = requests.post(token_url, data=api_credentials, timeout=60)
        response.raise_for_status()
        return response.json()['access_token']
    except requests.RequestException as e:
        logging.error(f"获取访问令牌失败: {e}")
        return None

@retry(stop=stop_after_attempt(5), wait=wait_fixed(60))
async def fetch_data_from_api(session, url, headers, table_name, cols_param):
    async with semaphore:
        try:
            await asyncio.sleep(REQUESTS_INTERVAL)  # 保证请求间隔
            async with session.get(url, headers=headers, timeout=60) as response:
                if response.status == 429:
                    # 处理限流
                    retry_after = int(response.headers.get('Retry-After', 60))
                    logging.info(f"触发限流,{retry_after}秒后重试")
                    await asyncio.sleep(retry_after)
                    raise Exception("达到限流阈值")
                response.raise_for_status()
                data = await response.json()
                total_pages = data.get('total_pages', 1)
                records = data.get('records', {})
                all_records = records.get(table_name, [])
                tasks = [
                    fetch_paginated_data(session, f"{url}&page={page}", headers, table_name)
                    for page in range(2, total_pages + 1)
                ]
                paginated_results = await asyncio.gather(*tasks)
                for result in paginated_results:
                    if result:
                        all_records.extend(result)
                return all_records
        except (aiohttp.ClientError, asyncio.TimeoutError) as e:
            logging.error(f"从API获取数据失败: {e}")
            raise

@retry(stop=stop_after_attempt(5), wait=wait_fixed(60))
async def fetch_paginated_data(session, url, headers, table_name):
    async with semaphore:
        try:
            await asyncio.sleep(REQUESTS_INTERVAL)  # 保证请求间隔
            async with session.get(url, headers=headers, timeout=60) as response:
                response.raise_for_status()
                data = await response.json()
                records = data.get('records', {}).get(table_name, [])
                return records
        except (aiohttp.ClientError, asyncio.TimeoutError) as e:
            logging.error(f"获取分页数据失败({url}): {e}")
            raise

def truncate_table(engine, table_name):
    try:
        with engine.connect() as connection:
            connection.execute(text(f'TRUNCATE TABLE fluxx.{table_name} RESTART IDENTITY CASCADE'))
            connection.commit()  # 提交事务
            logging.info(f"表 {table_name} 清空成功")
    except Exception as e:
        logging.error(f"清空表 {table_name} 失败: {e}")

def insert_data_to_postgres_sync(engine, table_name, column_names, data):
    df = pd.DataFrame(data, columns=column_names)

    if table_name == 'grant_request':
        current_utc_time = datetime.now(timezone.utc)
        df['sync_at'] = current_utc_time  # 为grant_request表添加同步时间

    try:
        # 插入数据到PostgreSQL
        df.to_sql(table_name, engine, schema='fluxx', if_exists='append', index=False)
        logging.info(f"数据插入表 {table_name} 成功")
    except Exception as e:
        logging.error(f"插入数据到表 {table_name} 失败: {e}")

async def process_table(session, url, headers, table_name, column_names, engine):
    try:
        data = await fetch_data_from_api(session, url, headers, table_name, column_names)
        if data is None:
            logging.error(f"获取表 {table_name} 数据失败,跳过")
            return
        # 插入前清空表
        truncate_table(engine, table_name)
        # 插入数据
        insert_data_to_postgres_sync(engine, table_name, column_names, data)
    except RetryError:
        logging.error(f"多次重试后仍无法获取表 {table_name} 数据,跳过")

async def extractor():
    api_credentials = {
        "grant_type": "xxxxxxxxxxxxxx",
        "client_id": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx",
        "client_secret": "xxxxxxxxxxxxxxxxxxxxxx"
    }

    access_token = await fetch_access_token(api_credentials)
    if not access_token:
        logging.error("未获取到访问令牌,退出")
        return

    headers = {
        "Authorization": f"Bearer {access_token}"
    }

    # 数据库配置
    db_config = {
        "database": "xxxxxxxxxxxxxxxxxx",
        "user": "xxxxxxxxxxxxxxx",
        "password": "xxxxxxxxxxxxxx",
        "host": "xxxxxxxxxxxxxxx",
        "port": "xxxxxxxxxxxxxxxxxx"
    }

    # 创建同步SQLAlchemy引擎
    db_url_sync = f"postgresql://{db_config['user']}:{db_config['password']}@{db_config['host']}:{db_config['port']}/{db_config['database']}"
    sync_engine = create_engine(db_url_sync)

    async with aiohttp.ClientSession() as session:
        async with asyncpg.create_pool(**db_config) as pool:
            async with pool.acquire() as conn:
                async with conn.transaction():
                    table_names_query = "SELECT table_name, array_agg(column_name) FROM column_metadata GROUP BY table_name"
                    rows = await conn.fetch(table_names_query)

            for row in rows:
                table_name, column_names = row['table_name'], row['array_agg']
                cols_param = ','.join([f'"{col}"' for col in column_names])
                url = f"https://xxxxxxxxxxxxxxxx/api/rest/v2/{table_name}?cols=[{cols_param}]&per_page={RECORDS_PER_PAGE}"
                await process_table(session, url, headers, table_name, column_names, sync_engine)

def data_extractor():
    asyncio.run(extractor())
    return "extraction is successfull"

local.settings.json配置

{
  "IsEncrypted": false,
  "Values": {
    "AzureWebJobsStorage": "UseDevelopmentStorage=true",
    "FUNCTIONS_WORKER_RUNTIME": "python",
    "AzureWebJobsFeatureFlags": "EnableWorkerIndexing"
  }
}

host.json配置

{
  "version": "2.0",
  "logging": {
    "applicationInsights": {
      "samplingSettings": {
        "isEnabled": true,
        "excludedTypes": "Request"
      }
    }
  },
  "extensionBundle": {
    "id": "Microsoft.Azure.Functions.ExtensionBundle",
    "version": "[3.*, 4.0.0)"
  }
}

requirements.txt依赖

azure-functions
azure-functions-durable
pandas
requests
sqlalchemy
psycopg2-binary
asyncio
aiohttp
asyncpg
tenacity

解决方案

针对部署后无法检测到HTTP触发器的问题,可按以下步骤排查修复:

1. 在云端启用Worker Indexing

本地配置中已设置AzureWebJobsFeatureFlags: EnableWorkerIndexing,需确保云端Function App的应用设置中添加同样的配置:

  • 进入Azure门户的Function App → 配置 → 应用设置
  • 添加新设置:名称AzureWebJobsFeatureFlags,值EnableWorkerIndexing
  • 保存后重启Function App

2. 检查extensionBundle版本兼容性

当前host.json中使用的extensionBundle版本是[3.*, 4.0.0),建议升级到最新的稳定版本(如[4.*, 5.0.0)),确保与Durable Functions扩展兼容:

"extensionBundle": {
  "id": "Microsoft.Azure.Functions.ExtensionBundle",
  "version": "[4.*, 5.0.0)"
}

3. 确认部署文件结构

Python Function App要求文件结构符合规范:

FunctionAppName/
├── host.json
├── requirements.txt
├── <FunctionName>/
│   ├── __init__.py  # 核心函数代码所在文件
└── fluxx_extraction.py

确保核心函数代码位于单独的文件夹(如HttpStart/)下的__init__.py中,而非根目录的__init__.py(如果是单函数项目,根目录的__init__.py也可,但需确保部署时所有文件都上传)。

4. 验证运行时设置

确认Function App的运行时设置正确:

  • 进入Function App → 配置 → 常规设置
  • 运行时堆栈选择Python,版本选择与本地开发一致的版本(如3.9/3.10)

5. 检查部署日志与依赖安装

  • 进入Function App → 部署中心 → 日志,查看部署过程中是否有依赖安装失败的报错
  • 如果psycopg2-binary安装失败,可尝试替换为psycopg2,或在部署前使用pip wheel预编译依赖包后再部署

6. 查看函数运行日志

  • 进入Function App → 监测 → 日志,查看是否有函数初始化时的报错,比如模块导入失败、配置缺失等问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 22:52:01