跳到主要内容

HTTP 节点

概述

HTTP 节点用于在离线开发工作流中发送 GET 或 POST 请求。响应状态码以 2 开头(如 200、201)时,节点运行成功;其他状态码会使节点运行失败。对于耗时较长的异步接口,还可以通过回调轮询持续查询任务状态。

节点配置

节点名称和描述

填写「节点名称」,并按需填写「描述」,便于在工作流画布和运行记录中识别节点。

请求地址

选择 GET 或 POST 请求方式,并填写接口 URL。默认请求方式为 GET。

参数配置

配置项说明
「请求头」以 key/value 形式配置请求头,例如资源类型等信息。
请求头
「认证」支持「API 密钥」和「Token 令牌」两种认证方式。API 密钥用于标识请求来源;Token 令牌在有效期内代表用户身份和接口访问权限。
「请求参数」配置接口需要的字段和查询条件。支持引用全局参数、时间宏参数和工作流参数。

「请求体」以 JSON 格式填写 POST 请求体。支持引用全局参数、时间宏参数和工作流参数。

回调轮询

回调轮询配置

当接口请求时间较长时,可配置回调接口轮询访问请求状态。

配置项说明
「轮询间隔时长(秒)」设置两次轮询之间的时间间隔。
「回调地址」选择 GET 或 POST 请求方式,并填写回调接口 URL。默认请求方式为 GET。
「请求头」以 key/value 形式配置请求头,例如资源类型等信息。
「认证」支持「API 密钥」和「Token 令牌」两种认证方式。
「请求参数」配置回调接口需要的字段和查询条件。支持引用全局参数、时间宏参数和工作流参数。
「请求体」以 JSON 格式填写 POST 请求体。支持引用全局参数、时间宏参数和工作流参数。

回调规则

  1. 主请求的响应必须包含 taskId,格式如下:
{
"taskId": "xxxxxx"
}
  1. 系统会自动提取 taskId,并将其拼接到回调地址末尾。

    例如,回调地址为 http://127.0.0.1:8004/status 时,拼接后的地址为 http://127.0.0.1:8004/status/xxxxxx。

  2. 回调响应必须包含 status,格式如下:

    {
    "status": "xxxxxx"
    }
  3. 回调响应状态码以 2 开头时,系统继续校验 status;其他状态码会使节点报错并结束运行。

  4. status 为 PROCESSING 时,系统继续轮询;为 SUCCESS 时,节点运行成功;为其他值时,节点运行失败。请确保回调服务最终返回结束状态,避免持续轮询。

运行选项

|500

  • 运行标志

    • 禁止执行:工作流运行至该节点后将直接跳过执行,常用于临时数据问题排查、部分任务运行控制等场景。
    • 正常:按照既有调度策略运行该节点,节点默认运行标志。
  • 失败重试

    • 重试次数:节点失败后的自动重试次数,默认1次
    • 重试间隔时间:触发重试的时间间隔,默认5分钟
  • 超时限制

    • 超时时间:单个节点的超时时间限制,超过后会自动失败

通过 HTTP 调用 Python 服务

以下示例通过 HTTP 节点调用 Python 服务,并使用回调轮询查询异步任务状态。

准备 Python API 服务所需依赖

flask
request
jsonify
cachetools

编写 Python 服务脚本

# -*- coding: utf-8 -*-

from flask import Flask, request, jsonify
from concurrent.futures import ThreadPoolExecutor
from cachetools import TTLCache
import uuid
import io
import contextlib


app = Flask(__name__)

# 线程池(最大 10 个并发任务)
executor = ThreadPoolExecutor(max_workers=10)

# 2 天转换为秒
two_days_in_seconds = 2 * 24 * 60 * 60

# 创建一个缓存对象,每个条目的 TTL 为 2 天
task_cache = TTLCache(maxsize=100, ttl=two_days_in_seconds)
def run_python_code(code, taskId):
""" 执行 Python 代码,并存储结果 """
try:
result = execute_code(code)
task_cache[taskId] = {"taskId": taskId, "status": "SUCCESS", "log": result}
except Exception as e:
task_cache[taskId] = {"taskId": taskId, "status": "FAILED", "error": str(e)}

@app.route("/run_code", methods=["POST"])
def run_code():
""" 提交代码任务,异步执行 """
data = request.get_json()
code = data.get("code", "").strip()

if not code:
return jsonify({"error": "No code provided"}), 400

# 生成唯一 taskId
taskId = str(uuid.uuid4())

# 记录任务开始
task_cache[taskId] = {"taskId": taskId, "status": "PROCESSING"}

# 提交任务到线程池
executor.submit(run_python_code, code, taskId)

return jsonify({"taskId": taskId, "status": "PROCESSING"})

@app.route("/status/<taskId>", methods=["GET"])
def check_status(taskId):
""" 查询任务状态 """
if taskId not in task_cache:
return jsonify({"taskId": taskId, "status": "FAILED", "error": "Task not found"}), 404

task_info = task_cache[taskId]

if task_info["status"] == "FAILED":
return jsonify({"taskId": taskId, "status": "FAILED", "error": task_info["error"]}), 500
return jsonify(task_info)

# 执行代码的函数
def execute_code(code):
# 使用 StringIO 捕获 print 输出
output_buffer = io.StringIO()
try:
# 重定向标准输出到 output_buffer
with contextlib.redirect_stdout(output_buffer):
exec(code)
# 获取捕获的输出
output = output_buffer.getvalue()
return output.strip() # 返回捕获的输出
except Exception as e:
# 如果代码执行失败,返回错误信息
return f"Code execution failed: {str(e)}"
finally:
# 关闭 StringIO 缓冲区
output_buffer.close()


if __name__ == "__main__":
app.run(host="0.0.0.0", port=8004, debug=True) # 监听所有地址,端口 8004

配置 HTTP 节点

  1. 将请求方式设置为 POST,并在「请求体」中配置参数。

  2. 开启「回调轮询」,并配置回调接口。