处理引擎和 Python 插件
使用 InfluxDB 3 Core 中的处理引擎,可通过自定义 Python 代码扩展您的数据库。您可以设置在写入时、定时或按需触发代码,以实现工作流自动化、数据转换并创建 API 端点。
什么是处理引擎?
处理引擎是一个嵌入在 InfluxDB 3 Core 数据库中的 Python 虚拟机。您可以配置触发器,以便在以下情况下运行您的 Python 插件代码:
- 数据写入 - 在数据进入数据库时对其进行处理和转换
- 定时事件 - 在定义的间隔或特定时间运行代码
- HTTP 请求 - 开放可执行您代码的自定义 API 端点
您可以使用处理引擎的内存缓存来管理执行间的状态,并直接在数据库中构建有状态应用程序。
本指南将引导您完成设置处理引擎、创建第一个插件以及配置在特定事件上执行代码的触发器的过程。
开始之前
确保您具备
- 一个可用的 InfluxDB 3 Core 实例
- 命令行访问权限
- 如果您要编写自己的插件,需安装 Python
- InfluxDB CLI 的基本知识
准备好上述先决条件后,请按照以下步骤实施处理引擎,以满足您的数据自动化需求。
设置处理引擎
当配置了 --plugin-dir 或 INFLUXDB3_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 时,请按照以下步骤配置处理引擎:
- 决定每个插件的运行位置
- 数据处理插件(如 WAL 插件)在摄取节点上运行
- HTTP 触发的插件在处理 API 请求的节点上运行
- 定时插件可以在任何已配置的节点上运行
- 在正确的实例上启用插件
- 确保所有运行插件的实例上插件文件保持一致
- 使用共享存储或文件同步工具来保持插件一致性
向运行插件的节点提供插件
在与运行触发器和插件的节点相同的系统上配置您的插件目录。
添加处理引擎插件
插件是一个 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"}下一步
编写完插件后
- 创建触发器以将插件连接到数据库事件
- 安装插件所需的所有 Python 依赖项
- 了解如何使用 API 扩展插件
从本地机器上传插件
对于本地开发和测试,您可以在创建触发器时直接从您的机器上传插件文件。这免去了手动将文件复制到服务器插件目录的步骤。
使用 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 请求时运行。
要运行插件,请发送 GET 或 POST 请求到该端点——例如:
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_monitorcurl -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, requests 或 influxdb3-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 3 Core 和本文档提供反馈和错误报告。要获得支持,请使用以下资源
具有年度合同或支持合同的客户可以 联系 InfluxData 支持。