文档文档

处理引擎和 Python 插件

使用 InfluxDB 3 Core 中的处理引擎,可通过自定义 Python 代码扩展您的数据库。您可以设置在写入时、定时或按需触发代码,以实现工作流自动化、数据转换并创建 API 端点。

什么是处理引擎?

处理引擎是一个嵌入在 InfluxDB 3 Core 数据库中的 Python 虚拟机。您可以配置触发器,以便在以下情况下运行您的 Python 插件代码:

  • 数据写入 - 在数据进入数据库时对其进行处理和转换
  • 定时事件 - 在定义的间隔或特定时间运行代码
  • HTTP 请求 - 开放可执行您代码的自定义 API 端点

您可以使用处理引擎的内存缓存来管理执行间的状态,并直接在数据库中构建有状态应用程序。

本指南将引导您完成设置处理引擎、创建第一个插件以及配置在特定事件上执行代码的触发器的过程。

开始之前

确保您具备

  • 一个可用的 InfluxDB 3 Core 实例
  • 命令行访问权限
  • 如果您要编写自己的插件,需安装 Python
  • InfluxDB CLI 的基本知识

准备好上述先决条件后,请按照以下步骤实施处理引擎,以满足您的数据自动化需求。

设置处理引擎

当配置了 --plugin-dirINFLUXDB3_PLUGIN_DIR 时,处理引擎将被激活。

不同部署类型的默认行为

部署方式默认状态配置
Docker 镜像已启用INFLUXDB3_PLUGIN_DIR=/plugins
DEB/RPM 包已启用plugin-dir="/var/lib/influxdb3/plugins"
二进制/源码已禁用未配置 plugin-dir

如果您使用 Docker 或 DEB/RPM 包安装了 InfluxDB 3 Core,处理引擎已启用——请跳至添加处理引擎插件。要禁用处理引擎,请参阅启用和禁用处理引擎

手动启用处理引擎

若要通过二进制文件或源代码构建运行并激活处理引擎,请使用 --plugin-dir 标志启动 InfluxDB 3 Core 服务器。该标志会告知 InfluxDB 从何处加载插件文件。

将 influxdb3 二进制文件与其 python 目录放在一起

influxdb3 二进制文件需要相邻的 python/ 目录才能正常运行。如果您手动从 tar.gz 解压,请将其保留在同一个父目录中。

your-install-location/
├── influxdb3
└── python/

将父目录添加到您的 PATH 中;不要将二进制文件移出此目录。

influxdb3 serve \
  --
NODE_ID
\
--object-store
OBJECT_STORE_TYPE
\
--plugin-dir
PLUGIN_DIR

在上面的示例中,替换以下内容

  • NODE_ID:实例的唯一标识符
  • OBJECT_STORE_TYPE:对象存储类型(例如 file 或 s3)
  • PLUGIN_DIR:插件文件存储目录的绝对路径。请将所有插件文件存储在此目录或其子目录中。

使用自定义插件仓库

默认情况下,以 gh: 为前缀引用的插件将从官方 influxdata/influxdb3_plugins 仓库中获取。要使用自定义仓库,请在启动服务器时添加 --plugin-repo 标志。详情请参阅使用自定义插件仓库

配置分布式环境

在分布式设置中运行 InfluxDB 3 Core 时,请按照以下步骤配置处理引擎:

  1. 决定每个插件的运行位置
    • 数据处理插件(如 WAL 插件)在摄取节点上运行
    • HTTP 触发的插件在处理 API 请求的节点上运行
    • 定时插件可以在任何已配置的节点上运行
  2. 在正确的实例上启用插件
  3. 确保所有运行插件的实例上插件文件保持一致
    • 使用共享存储或文件同步工具来保持插件一致性

向运行插件的节点提供插件

在与运行触发器和插件的节点相同的系统上配置您的插件目录。

添加处理引擎插件

插件是一个 Python 脚本,它定义了一个具有兼容触发器(触发器规范)签名的函数。当指定的事件发生时,InfluxDB 会运行该插件。

选择插件策略

您有两种主要方式为 InfluxDB 实例添加插件:

使用示例插件

InfluxData 维护着一个官方和社区插件的仓库,您可以立即在处理引擎设置中使用它们。

浏览插件库以查找示例及 InfluxData 官方插件,用于:

  • 数据转换:处理和转换传入的数据
  • 警报:根据数据阈值发送通知
  • 聚合:对时间序列数据计算统计信息
  • 集成:连接到外部服务和 API
  • 系统监控:跟踪资源使用情况和健康指标

有关社区贡献的插件,请查看 GitHub 上的 influxdb3_plugins 仓库

添加示例插件

使用仓库中插件有两种选择:

选项 1:本地复制插件

克隆 influxdata/influxdb3_plugins 仓库并将插件复制到您配置的插件目录中。

# Clone the repository
git clone https://github.com/influxdata/influxdb3_plugins.git
   
# Copy a plugin to your configured plugin directory
cp influxdb3_plugins/influxdata/system_metrics/system_metrics.py /path/to/plugins/
选项 2:直接从 GitHub 引用插件

通过使用 gh: 前缀直接从 GitHub 引用插件,跳过下载步骤。

# Create a trigger using a plugin from GitHub
influxdb3 create trigger \
  --trigger-spec "every:1m" \
  --path "gh:influxdata/system_metrics/system_metrics.py" \
  --database my_database \
  system_metrics

这种方法:

  • 确保您使用的是最新版本
  • 简化更新和维护
  • 减少本地存储需求
选项 3:使用自定义插件仓库

对于维护自己插件仓库或需要使用私有/内部插件的组织,请配置自定义插件仓库 URL。

# Start the server with a custom plugin repository
influxdb3 serve \
  --node-id node0 \
  --object-store file \
  --data-dir ~/.influxdb3 \
  --plugin-dir ~/.plugins \
  --plugin-repo "https://internal.company.com/influxdb-plugins/"

然后使用 gh: 前缀引用自定义仓库中的插件。

# Fetches from: https://internal.company.com/influxdb-plugins/myorg/custom_plugin.py
influxdb3 create trigger \
  --trigger-spec "every:5m" \
  --path "gh:myorg/custom_plugin.py" \
  --database my_database \
  custom_trigger

自定义仓库的使用场景

  • 私有插件:托管不适合公共仓库的专有插件
  • 物理隔离环境(Air-gapped):在无法访问外网时使用内部镜像
  • 开发和阶段性环境:在部署到生产环境前测试开发分支中的插件
  • 合规要求:满足要求内部托管的数据治理政策

--plugin-repo 选项接受任何提供原始插件文件的 HTTP/HTTPS URL。详情请参阅 plugin-repo 配置选项

插件具有多种功能,例如:

  • 接收插件特定的参数(如写入的数据、调用时间或 HTTP 请求)
  • 访问通过触发器参数配置传入的关键字参数(即 args
  • 访问 influxdb3_local 共享 API 以写入数据、查询数据及在执行间管理状态

有关可用函数、参数以及插件如何与 InfluxDB 交互的更多信息,请参阅如何扩展插件

创建自定义插件

要构建自定义功能,您可以创建自己的处理引擎插件。

前提条件

开始之前,请确保:

  • 处理引擎已在您的 InfluxDB 3 Core 实例上启用。
  • 您已配置好存储插件文件的 --plugin-dir
  • 您有权访问该插件目录。

创建插件的步骤

选择插件类型

根据您的自动化目标选择插件类型:

插件类型适用场景
数据写入处理实时传入的数据
定时任务在指定间隔或时间运行代码
HTTP 请求通过 API 端点按需运行代码

创建插件文件

插件现在支持单文件和多文件架构

单文件插件

  • 在您的插件目录中创建一个 .py 文件
  • 根据所选插件类型添加相应的函数签名
  • 在函数内部编写处理逻辑

多文件插件

  • 在您的插件目录中创建一个目录
  • 添加一个 __init__.py 文件作为入口点(必需)
  • 将辅助模块组织在其他 .py 文件中
  • 在插件代码中导入并使用模块
多文件插件示例结构
my_plugin/
├── __init__.py       # Required - entry point with trigger function
├── utils.py          # Supporting module
├── processors.py     # Data processing functions
└── config.py         # Configuration helpers

__init__.py 文件必须包含您的触发函数

# my_plugin/__init__.py
from .processors import process_data
from .config import get_settings

def process_writes(influxdb3_local, table_batches, args=None):
    settings = get_settings()
    for table_batch in table_batches:
        process_data(influxdb3_local, table_batch, settings)

辅助模块可以包含辅助函数

# my_plugin/processors.py
def process_data(influxdb3_local, table_batch, settings):
    # Processing logic here
    pass

编写完插件后,创建触发器将其连接到数据库事件并定义何时运行。

创建数据写入插件

使用数据写入插件在数据写入数据库时进行处理。这些插件使用 table:all_tables: 触发器规范。理想的用例包括:

  • 数据转换和丰富
  • 对传入数值进行预警
  • 创建派生指标
def process_writes(influxdb3_local, table_batches, args=None):
    # Process data as it's written to the database
    for table_batch in table_batches:
        table_name = table_batch["table_name"]
        rows = table_batch["rows"]
        
        # Log information about the write
        influxdb3_local.info(f"Processing {len(rows)} rows from {table_name}")
        
        # Write derived data back to the database
        line = LineBuilder("processed_data")
        line.tag("source_table", table_name)
        line.int64_field("row_count", len(rows))
        influxdb3_local.write(line)

创建定时插件

定时插件使用 every:cron: 触发器规范在定义的间隔运行。适用于:

  • 周期性数据聚合
  • 报表生成
  • 系统健康检查
def process_scheduled_call(influxdb3_local, call_time, args=None):
    # Run code on a schedule
    
    # Query recent data
    results = influxdb3_local.query("SELECT * FROM metrics WHERE time > now() - INTERVAL '1 hour'")
    
    # Process the results
    if results:
        influxdb3_local.info(f"Found {len(results)} recent metrics")
    else:
        influxdb3_local.warn("No recent metrics found")

创建 HTTP 请求插件

HTTP 请求插件使用 request: 触发器规范响应 API 调用。适用于:

  • 创建自定义 API 端点
  • 外部集成的 Webhook
  • 用于数据交互的用户界面
def process_request(influxdb3_local, query_parameters, request_headers, request_body, args=None):
    # Handle HTTP requests to a custom endpoint
    
    # Log the request parameters
    influxdb3_local.info(f"Received request with parameters: {query_parameters}")
    
    # Process the request body
    if request_body:
        import json
        data = json.loads(request_body)
        influxdb3_local.info(f"Request data: {data}")
    
    # Return a response (automatically converted to JSON)
    return {"status": "success", "message": "Request processed"}

下一步

编写完插件后

从本地机器上传插件

对于本地开发和测试,您可以在创建触发器时直接从您的机器上传插件文件。这免去了手动将文件复制到服务器插件目录的步骤。

使用 influxdb3 CLI 上传插件

配合 --path 使用 --upload 标志来传输本地文件或目录

# Upload single-file plugin
influxdb3 create trigger \
  --trigger-spec "every:10s" \
  --path "/local/path/to/plugin.py" \
  --upload \
  --database metrics \
  my_trigger

# Upload multifile plugin directory
influxdb3 create trigger \
  --trigger-spec "every:30s" \
  --path "/local/path/to/plugin-dir" \
  --upload \
  --database metrics \
  complex_trigger

有关更多信息,请参阅 influxdb3 create trigger CLI 参考

使用 HTTP API 上传插件

要使用 HTTP API 上传插件文件,请发送 PUT 请求至 /api/v3/plugins/files 端点

PUT localhost:8181/api/v3/plugins/files

在您的请求中包含以下内容

  • Headers:
    • Authorization: Bearer 需携带您的管理员令牌
    • Content-Type: application/octet-stream
  • 查询参数:
    • path (字符串,必需):插件文件相对于插件目录的路径
# Upload a single-file plugin
curl -X PUT "localhost:8181/api/v3/plugins/files?path=plugin.py" \
  --header "Authorization: Bearer 
AUTH_TOKEN
"
\
--header "Content-Type: application/octet-stream" \ --data-binary "@/local/path/to/plugin.py"

替换 AUTH_TOKEN:您的 管理员令牌

需要管理员权限

插件上传需要管理员令牌。此安全措施可防止未经授权的代码在服务器上执行。

何时使用插件上传

  • 本地插件开发和测试
  • 在无法通过 SSH 访问服务器的情况下部署插件
  • 快速迭代插件代码
  • 在 CI/CD 流水线中实现插件部署自动化

更新现有插件

修改运行中的触发器的插件代码,无需重新创建触发器。这样您可以在保留触发器配置和历史记录的同时迭代插件开发。

使用 influxdb3 CLI 更新插件

使用 influxdb3 update trigger 命令

# Update single-file plugin
influxdb3 update trigger \
  --database metrics \
  --trigger-name my_trigger \
  --path "/path/to/updated/plugin.py"

# Update multifile plugin
influxdb3 update trigger \
  --database metrics \
  --trigger-name complex_trigger \
  --path "/path/to/updated/plugin-dir"

有关完整参考,请参见 influxdb3 update trigger

使用 HTTP API 更新插件

要使用 HTTP API 更新插件文件,请发送 PUT 请求至 /api/v3/plugins/files 端点

PUT localhost:8181/api/v3/plugins/files

在您的请求中包含以下内容

  • Headers:
    • Authorization: Bearer 需携带您的管理员令牌
    • Content-Type: application/octet-stream
  • 查询参数:
    • path (字符串,必需):插件文件相对于插件目录的路径
# Update a plugin file
curl -X PUT "localhost:8181/api/v3/plugins/files?path=plugin.py" \
  --header "Authorization: Bearer 
AUTH_TOKEN
"
\
--header "Content-Type: application/octet-stream" \ --data-binary "@/path/to/updated/plugin.py"

替换 AUTH_TOKEN:您的 管理员令牌

更新操作

  • 立即替换插件文件
  • 保留触发器配置(规范、时间表、参数)
  • 出于安全考虑,需要管理员令牌
  • 适用于本地路径和上传的文件

查看已加载的插件

监控系统中加载了哪些插件,以实现运维透明度和故障排除。

选项 1:使用 CLI 命令

# List all plugins
influxdb3 show plugins --token $ADMIN_TOKEN

# JSON format for programmatic access
influxdb3 show plugins --format json --token $ADMIN_TOKEN

选项 2:查询系统表

_internal 数据库中的 system.plugin_files 表提供了详细的插件文件信息

influxdb3 query \
  -d _internal \
  "SELECT * FROM system.plugin_files ORDER BY plugin_name" \
  --token $ADMIN_TOKEN

可用列

  • plugin_name (字符串):触发器名称
  • file_name (字符串):插件文件名
  • file_path (字符串):服务器完整路径
  • size_bytes (Int64):文件大小
  • last_modified (Int64):修改时间戳(毫秒)

示例查询

-- Find plugins by name
SELECT * FROM system.plugin_files WHERE plugin_name = 'my_trigger';

-- Find large plugins
SELECT plugin_name, size_bytes
FROM system.plugin_files
WHERE size_bytes > 10000;

-- Check modification times
SELECT plugin_name, file_name, last_modified
FROM system.plugin_files
ORDER BY last_modified DESC;

更多信息,请参阅 influxdb3 show plugins 参考查询系统数据

创建触发器

触发器将您的插件代码连接到数据库事件。当指定事件发生时,处理引擎会执行您的插件。

了解触发器类型

插件类型触发器规范插件运行时间
数据写入table:<TABLE_NAME>all_tables数据写入到表时
定时任务every:<DURATION>cron:<EXPRESSION>在指定的时间间隔
HTTP 请求request:<REQUEST_PATH>接收到 HTTP 请求时

使用 influxdb3 CLI 创建触发器

使用带有相应触发器规范的 influxdb3 create trigger 命令

influxdb3 create trigger \
  --trigger-spec 
SPECIFICATION
\
--path
PLUGIN_FILE
\
--database
DATABASE_NAME
\
TRIGGER_NAME

在上面的示例中,替换以下内容

  • SPECIFICATION:触发器规范
  • PLUGIN_FILE:相对于您配置的插件目录的插件文件名
  • DATABASE_NAME:数据库名称
  • TRIGGER_NAME:新触发器的名称

插件路径

  • 对于单文件插件,只需向 --path 提供 .py 文件名(例如 test_plugin.py)。
  • 对于多文件插件,提供包含 __init__.py 的目录名称。

当不使用 --upload 时,服务器将相对于配置的 --plugin-dir 解析路径。有关多文件插件结构的详细信息,请参阅 创建您的插件文件

完整参考请参阅 influxdb3 create trigger

使用 HTTP API 创建触发器

要使用 HTTP API 创建触发器,请发送 POST 请求至 /api/v3/configure/processing_engine_trigger 端点

POST localhost:8181/api/v3/configure/processing_engine_trigger

在您的请求中包含以下内容

  • Headers:
    • Authorization: Bearer 加上您的身份验证令牌
    • Content-Type: application/json
  • 请求体:包含触发器配置的 JSON 对象
    • db (字符串,必填):数据库名称
    • trigger_name (字符串,必需):触发器名称
    • plugin_filename (字符串,必需):相对于插件目录的插件文件名
    • trigger_specification (字符串,必需):插件运行的时间(参见 触发器类型
    • trigger_settings (对象,必需):错误处理和执行的配置
      • run_async (布尔值):是否异步运行(默认:false
      • error_behavior (字符串):如何处理错误:Log(记录)、Retry(重试)或 Disable(禁用)(默认:Log
    • disabled (布尔值,必需):触发器是否已禁用
    • trigger_arguments (对象,可选):传递给插件的参数
# Create a basic trigger
curl -X POST "localhost:8181/api/v3/configure/processing_engine_trigger" \
  --header "Authorization: Bearer 
AUTH_TOKEN
"
\
--header "Content-Type: application/json" \ --data '{ "db": "
DATABASE_NAME
",
"trigger_name": "
TRIGGER_NAME
",
"plugin_filename": "
PLUGIN_FILE
",
"trigger_specification": "
TRIGGER_SPEC
",
"trigger_settings": { "run_async": false, "error_behavior": "Log" }, "disabled": false }'

在上面的示例中,替换以下内容

  • DATABASE_NAME:数据库名称
  • TRIGGER_NAME:新触发器的名称
  • PLUGIN_FILE:相对于您配置的插件目录的插件文件名
  • TRIGGER_SPEC:触发器规范(参见 示例
  • AUTH_TOKEN:您的令牌

触发器规范示例

以下示例演示如何为不同事件类型创建触发器。

在数据写入时触发

# Trigger on writes to a specific table
# The plugin file must be in your configured plugin directory
influxdb3 create trigger \
  --trigger-spec "table:sensor_data" \
  --path "process_sensors.py" \
  --database my_database \
  sensor_processor

# Trigger on writes to all tables
influxdb3 create trigger \
  --trigger-spec "all_tables" \
  --path "process_all_data.py" \
  --database my_database \
  all_data_processor
# Trigger on writes to a specific table
curl -X POST "localhost:8181/api/v3/configure/processing_engine_trigger" \
  --header "Authorization: Bearer 
AUTH_TOKEN
"
\
--header "Content-Type: application/json" \ --data '{ "db": "
DATABASE_NAME
",
"trigger_name": "sensor_processor", "plugin_filename": "process_sensors.py", "trigger_specification": "table:sensor_data", "trigger_settings": { "run_async": false, "error_behavior": "Log" }, "disabled": false }' # Trigger on writes to all tables curl -X POST "localhost:8181/api/v3/configure/processing_engine_trigger" \ --header "Authorization: Bearer
AUTH_TOKEN
"
\
--header "Content-Type: application/json" \ --data '{ "db": "
DATABASE_NAME
",
"trigger_name": "all_data_processor", "plugin_filename": "process_all_data.py", "trigger_specification": "all_tables", "trigger_settings": { "run_async": false, "error_behavior": "Log" }, "disabled": false }'

替换以下内容:

  • DATABASE_NAME:数据库名称
  • AUTH_TOKEN:您的令牌

当数据库将指定表的摄取数据刷新到对象存储中的预写日志 (WAL) 时,触发器就会运行(默认每秒一次)。

插件接收写入的数据和表信息。

带表排除的数据写入触发器

如果您希望对所有表使用单个触发器,但需要排除特定表,则可以使用触发器参数并在插件代码中过滤掉不需要的表——例如:

influxdb3 create trigger \
  --database 
DATABASE_NAME
\
--token
AUTH_TOKEN
\
--path processor.py \ --trigger-spec "all_tables" \ --trigger-arguments "exclude_tables=temp_data,debug_info,system_logs" \ data_processor

替换以下内容:

  • DATABASE_NAME:数据库名称
  • AUTH_TOKEN:您的 令牌

然后,在您的插件中:

# processor.py
def on_write(self, database, table_name, batch):
    # Get excluded tables from trigger arguments
    excluded_tables = set(self.args.get('exclude_tables', '').split(','))

    if table_name in excluded_tables:
        return

    # Process allowed tables
    self.process_data(database, table_name, batch)
建议
  • 提前返回:在插件中尽可能早地检查排除项。
  • 高效查询:对于较大的排除列表,使用集合以获得 O(1) 的查询性能。
  • 性能:记录跳过的表以供调试,但避免在生产环境中进行过多的日志记录。
  • 多个触发器:对于少量表,考虑创建多个特定的表触发器,而不是在插件代码内过滤。有关管理触发器的详情,请参阅 HTTP API 处理引擎端点

定时触发

# Run every 5 minutes
influxdb3 create trigger \
  --trigger-spec "every:5m" \
  --path "periodic_check.py" \
  --database my_database \
  regular_check

# Run on a cron schedule (8am daily)
# Supports extended cron format with seconds
influxdb3 create trigger \
  --trigger-spec "cron:0 0 8 * * *" \
  --path "daily_report.py" \
  --database my_database \
  daily_report
# Run every 5 minutes
curl -X POST "localhost:8181/api/v3/configure/processing_engine_trigger" \
  --header "Authorization: Bearer 
AUTH_TOKEN
"
\
--header "Content-Type: application/json" \ --data '{ "db": "
DATABASE_NAME
",
"trigger_name": "regular_check", "plugin_filename": "periodic_check.py", "trigger_specification": "every:5m", "trigger_settings": { "run_async": false, "error_behavior": "Log" }, "disabled": false }' # Run on a cron schedule (8am daily) # Supports extended cron format with seconds curl -X POST "localhost:8181/api/v3/configure/processing_engine_trigger" \ --header "Authorization: Bearer
AUTH_TOKEN
"
\
--header "Content-Type: application/json" \ --data '{ "db": "
DATABASE_NAME
",
"trigger_name": "daily_report", "plugin_filename": "daily_report.py", "trigger_specification": "cron:0 0 8 * * *", "trigger_settings": { "run_async": false, "error_behavior": "Log" }, "disabled": false }'

替换以下内容:

  • DATABASE_NAME:数据库名称
  • AUTH_TOKEN:您的令牌

插件接收定时调用的时间。

HTTP 请求触发

# Create an endpoint at /api/v3/engine/webhook
influxdb3 create trigger \
  --trigger-spec "request:webhook" \
  --path "webhook_handler.py" \
  --database my_database \
  webhook_processor
# Create an endpoint at /api/v3/engine/webhook
curl -X POST "localhost:8181/api/v3/configure/processing_engine_trigger" \
  --header "Authorization: Bearer 
AUTH_TOKEN
"
\
--header "Content-Type: application/json" \ --data '{ "db": "
DATABASE_NAME
",
"trigger_name": "webhook_processor", "plugin_filename": "webhook_handler.py", "trigger_specification": "request:webhook", "trigger_settings": { "run_async": false, "error_behavior": "Log" }, "disabled": false }'

替换以下内容:

  • DATABASE_NAME:数据库名称
  • AUTH_TOKEN:您的令牌

/api/v3/engine/{REQUEST_PATH} 访问您的端点(在此示例中为 /api/v3/engine/webhook)。触发器默认启用,并在指定路径收到 HTTP 请求时运行。

要运行插件,请发送 GETPOST 请求到该端点——例如:

curl https://:8181/api/v3/engine/webhook

插件接收带有方法、标头和正文的 HTTP 请求对象。

要查看与数据库关联的触发器,请使用 influxdb3 show summary 命令

influxdb3 show summary --database my_database --token AUTH_TOKEN

将参数传递给插件

使用触发器参数将配置从触发器传递给它运行的插件。您可以将其用于:

  • 监控阈值
  • 外部服务的连接属性
  • 插件行为的配置设置
influxdb3 create trigger \
  --trigger-spec "every:1h" \
  --path "threshold_check.py" \
  --trigger-arguments threshold=90,notify_email=admin@example.com \
  --database my_database \
  threshold_monitor
curl -X POST "localhost:8181/api/v3/configure/processing_engine_trigger" \
  --header "Authorization: Bearer 
AUTH_TOKEN
"
\
--header "Content-Type: application/json" \ --data '{ "db": "
DATABASE_NAME
",
"trigger_name": "threshold_monitor", "plugin_filename": "threshold_check.py", "trigger_specification": "every:1h", "trigger_settings": { "run_async": false, "error_behavior": "Log" }, "trigger_arguments": { "threshold": "90", "notify_email": "admin@example.com" }, "disabled": false }'

替换以下内容:

  • DATABASE_NAME:数据库名称
  • AUTH_TOKEN:您的令牌

参数作为 Dict[str, str] 传递给插件,其中键是参数名称,值是参数值

def process_scheduled_call(influxdb3_local, call_time, args=None):
    if args and "threshold" in args:
        threshold = float(args["threshold"])
        email = args.get("notify_email", "default@example.com")
        
        # Use the arguments in your logic
        influxdb3_local.info(f"Checking threshold {threshold}, will notify {email}")

控制触发器执行

默认情况下,触发器同步运行——每个实例都会等待前一个实例完成后再执行。

要允许同一触发器的多个实例同时运行,请配置触发器异步运行

# Allow multiple trigger instances to run simultaneously
influxdb3 create trigger \
  --trigger-spec "table:metrics" \
  --path "heavy_process.py" \
  --run-asynchronous \
  --database my_database \
  async_processor
# Allow multiple trigger instances to run simultaneously
curl -X POST "localhost:8181/api/v3/configure/processing_engine_trigger" \
  --header "Authorization: Bearer 
AUTH_TOKEN
"
\
--header "Content-Type: application/json" \ --data '{ "db": "
DATABASE_NAME
",
"trigger_name": "async_processor", "plugin_filename": "heavy_process.py", "trigger_specification": "table:metrics", "trigger_settings": { "run_async": true, "error_behavior": "Log" }, "disabled": false }'

替换以下内容:

  • DATABASE_NAME:数据库名称
  • AUTH_TOKEN:您的令牌

配置触发器的错误处理

要配置触发器的错误处理行为,请指定以下值之一:

  • log(默认):将所有插件错误记录到 stdout 和触发器数据库中的 system.processing_engine_logs 表。
  • retry:在错误发生后立即尝试再次运行插件。
  • disable:在错误发生时自动禁用插件(以后可以重新启用)。

有关更多信息,请参见如何查询触发器日志

# Automatically retry on error
influxdb3 create trigger \
  --trigger-spec "table:important_data" \
  --path "critical_process.py" \
  --error-behavior retry \
  --database my_database \
  critical_processor

# Disable the trigger on error
influxdb3 create trigger \
  --trigger-spec "request:webhook" \
  --path "webhook_handler.py" \
  --error-behavior disable \
  --database my_database \
  auto_disable_processor
# Automatically retry on error
curl -X POST "localhost:8181/api/v3/configure/processing_engine_trigger" \
  --header "Authorization: Bearer 
AUTH_TOKEN
"
\
--header "Content-Type: application/json" \ --data '{ "db": "
DATABASE_NAME
",
"trigger_name": "critical_processor", "plugin_filename": "critical_process.py", "trigger_specification": "table:important_data", "trigger_settings": { "run_async": false, "error_behavior": "Retry" }, "disabled": false }' # Disable the trigger on error curl -X POST "localhost:8181/api/v3/configure/processing_engine_trigger" \ --header "Authorization: Bearer
AUTH_TOKEN
"
\
--header "Content-Type: application/json" \ --data '{ "db": "
DATABASE_NAME
",
"trigger_name": "auto_disable_processor", "plugin_filename": "webhook_handler.py", "trigger_specification": "request:webhook", "trigger_settings": { "run_async": false, "error_behavior": "Disable" }, "disabled": false }'

替换以下内容:

  • DATABASE_NAME:数据库名称
  • AUTH_TOKEN:您的令牌

管理插件依赖项

使用 influxdb3 install package 命令将第三方库(如 pandas, requestsinfluxdb3-python)添加到您的插件环境中。
这会将包安装到处理引擎的嵌入式 Python 环境中,以确保与您的 InfluxDB 实例的兼容性。

# Use the CLI to install a Python package
influxdb3 install package pandas
# Use the CLI to install a Python package in a Docker container
docker exec -it 
CONTAINER_NAME
influxdb3 install package pandas
# Use the HTTP API to install Python packages
curl -X POST "localhost:8181/api/v3/configure/plugin_environment/install_packages" \
  --header "Authorization: Bearer 
AUTH_TOKEN
"
\
--header "Content-Type: application/json" \ --data '{ "packages": ["pandas", "requests", "numpy"] }'

替换 AUTH_TOKEN:您的 管理员令牌

有关完整参考,请参阅 安装插件包

这些示例将指定的 Python 包(例如 pandas)安装到处理引擎的嵌入式虚拟环境中。

  • 当直接在系统上运行 InfluxDB 时,请使用 CLI 命令。
  • 如果您在容器化环境中运行 InfluxDB,请使用 Docker 变体。
  • 使用 HTTP API 进行程序化包安装或 CI/CD 工作流。

为插件使用捆绑的 Python

当您使用 --plugin-dir 选项启动服务器时,InfluxDB 3 会为您的插件创建一个 Python 虚拟环境 (<PLUGIN_DIR>/venv)。如果您需要创建自定义虚拟环境,请使用 InfluxDB 3 捆绑的 Python 解释器。不要使用系统 Python。使用系统 Python 创建虚拟环境(例如,使用 python -m venv)可能会导致运行时错误和插件失败。

更多信息,请参见 处理引擎 README

InfluxDB 会在您的插件目录中创建一个 Python 虚拟环境,并安装指定的包。

为安全环境禁用包安装

对于物理隔离部署或有严格安全要求的环境,您可以在保持处理引擎功能的同时禁用 Python 包安装。

使用 --package-manager disabled 启动服务器

influxdb3 serve \
  --node-id node0 \
  --object-store file \
  --data-dir ~/.influxdb3 \
  --plugin-dir ~/.plugins \
  --package-manager disabled

当包安装被禁用时:

  • 处理引擎继续正常处理触发器
  • 插件代码执行不受限制
  • 包安装命令被拦截
  • 虚拟环境中预安装的依赖项保持可用

预安装所需的依赖项

在禁用包管理器之前,请安装所有必需的 Python 包

# Install packages first
influxdb3 install package pandas requests numpy

# Then start with disabled package manager
influxdb3 serve \
  --plugin-dir ~/.plugins \
  --package-manager disabled

禁用包管理的使用场景

  • 无法访问互联网的物理隔离环境
  • 禁止运行时安装包的合规要求
  • 集中管理的依赖项环境
  • 要求仅允许预先批准的包的安全策略

更多配置选项,请参见 –package-manager

插件安全性

处理引擎包含安全功能,以保护您的 InfluxDB 3 Core 实例免受未经授权的代码执行和文件系统攻击。

插件路径验证

所有插件文件路径均经过验证,以防止目录遍历攻击。系统会拦截:

  • 带有父目录引用的相对路径 (../, ../../)
  • 绝对路径 (/etc/passwd, /usr/bin/script.py)
  • 逃逸出插件目录的符号链接

在创建或更新触发器时,插件路径必须解析在已配置的 --plugin-dir 内。

被阻止的路径示例

# These will be rejected
influxdb3 create trigger \
  --path "../../../etc/passwd" \  # Blocked: parent directory traversal
  ...

influxdb3 create trigger \
  --path "/tmp/malicious.py" \    # Blocked: absolute path
  ...

有效的插件路径

# These are allowed
influxdb3 create trigger \
  --path "myapp/plugin.py" \      # Relative to plugin-dir
  ...

influxdb3 create trigger \
  --path "transforms/data.py" \    # Subdirectory in plugin-dir
  ...

上传和更新权限

插件上传和更新操作需要管理员令牌,以防止未经授权的代码部署

  • --upload 标志需要管理员权限
  • update trigger 命令需要管理员令牌
  • 标准资源令牌无法上传或修改插件代码

此安全模型确保只有管理员才能在您的数据库中引入或修改可执行代码。

最佳实践

对于开发

  • 在开发期间使用 --upload 标志部署插件
  • 首先在非生产环境中测试插件
  • 在部署前审查插件代码

对于生产环境

  • 通过安全文件传输预先将插件部署到服务器的插件目录
  • 使用自定义插件仓库以获取经过审核和批准的插件
  • 在受锁定的环境中禁用包安装 (--package-manager disabled)
  • 使用 system.plugin_files 审计插件文件
  • 实施插件更新的变更控制流程

更多安全配置选项,请参见 配置选项


此页面是否有帮助?

感谢您的反馈!


InfluxDB OSS 2.9.0:API 令牌默认进行哈希处理

InfluxDB OSS 2.9.0 增强了令牌安全性 —— 令牌在磁盘上默认进行哈希处理。现有令牌在首次启动时会被哈希,之后无法恢复。请在升级前保存您仍然需要的所有明文令牌。

查看 InfluxDB OSS 2.9.0 发行说明

哈希令牌的认证方式与未哈希令牌完全相同 —— 客户端和集成功能可继续正常工作。

2.9.0 中的其他新特性

  • 可配置的备份压缩
  • 恢复对包含哈希令牌的备份的支持
  • 更严格的边缘数据复制(Edge Data Replication)队列验证
  • Flux 升级
  • 压缩可靠性改进

Explorer 1.9 的主要增强功能

Explorer 1.9 现已发布,支持 InfluxQL、AI 辅助的 Flux 转 SQL 转换器(测试版)以及新的实时示例数据模拟器。

查看 Explorer 1.9 发行说明

Explorer 1.9 包含多项新功能和改进,使查询、可视化和管理数据变得更加轻松。

亮点

  • Flux 转 SQL 转换器(测试版):通过 AI 辅助转换器将 Flux 查询转换为 SQL。
  • InfluxQL 支持:在数据浏览器(Data Explorer)和仪表板中使用 InfluxQL 查询数据,并保存和加载 InfluxQL 查询。
  • InfluxQL 可视化:根据 InfluxQL 结果渲染折线图和柱状图,并支持按标签进行序列分组。
  • 查询错误历史记录:在查询工具中查看查询错误历史记录。
  • 实时示例数据模拟器:使用新的鸟类数据和信号发生器模拟器生成连续的实时示例数据。

更多详细信息,请参阅 Explorer 1.9 发行说明

InfluxDB 3.10 现已发布

InfluxDB 3 Core 3.10 增加了自动目录格式升级、可配置的查询并发限制以及处理引擎改进。

InfluxDB 3 Core 3.10 的关键更新

  • 目录格式升级:在 3.10 首次启动时,磁盘目录会自动从 v2 格式升级到 v3 格式。迁移是单向的——升级前请务必备份您的目录。
  • --max-concurrent-queries:限制并发查询(可在运行时调整)。
  • GET /ready 端点,用于就绪探针。
  • 处理引擎:跨数据库查询和触发器锁定标志。

更多信息,请参阅 InfluxDB 3 Core 发行说明

InfluxDB 3.10 现已发布

InfluxDB 3 Enterprise 3.10 增加了自动备份与恢复、行级删除和用户管理功能,并改进了自动目录格式升级和性能预览。

InfluxDB 3 Enterprise 3.10 的关键更新

  • 目录格式升级:在 3.10 首次启动时,磁盘目录会自动从 v2 格式升级到 v3 格式。迁移是单向的——升级前请务必备份您的目录。
  • 自动备份与恢复(测试版)
  • 行级删除
  • 用户管理(认证和 RBAC)— 预览版
  • 性能预览改进

备份与恢复、行级删除以及性能预览需要升级到企业级存储引擎(选择性加入测试版)。测试版和预览版功能可能会发生重大变更,不建议用于生产环境。

更多信息,请参阅 InfluxDB 3 Enterprise 发行说明

Telegraf Enterprise 现已全面上市(General Availability)

Telegraf Enterprise 现已全面上市,同时发布的还有 Telegraf Controller v1.0

Telegraf Enterprise 将 Telegraf Controller(一个用于 Telegraf 的集中式管理控制台)与 InfluxData 的官方支持相结合。通过单一系统管理配置、监控集群健康状况并操作数以万计的 Telegraf 代理。

InfluxDB Docker 的 latest 标签将指向 InfluxDB 3 Core

2026 年 9 月 15 日起,InfluxDB Docker 镜像的 latest 标签将指向 InfluxDB 3 Core。为避免意外升级,请在 Docker 部署中使用特定版本标签。

如果使用 Docker 来安装和运行 InfluxDB,latest 标签将指向 InfluxDB 3 Core。为避免意外升级,请在您的 Docker 部署中使用特定的版本标签。例如,如果使用 Docker 运行 InfluxDB v2,请将 latest 版本标签替换为 Docker pull 命令中的特定版本标签 — 例如

docker pull influxdb:2