MinIO 桶通知实战:把对象事件发布到 Kafka、RabbitMQ、Redis 等十类目标端
本文以 MinIO 官方文档 docs/bucket/notifications/README.md 为主体,完整覆盖桶事件通知(Bucket Notification)的全部核心内容:支持的事件类型清单、十类通知目标(AMQP、MQTT、Elasticsearch、Redis、NATS、PostgreSQL、MySQL、Kafka、Webhook、NSQ)的完整配置参数与环境变量、mc 配置与验证命令、持久化事件存储(queue_dir/queue_limit)、namespace/access 双格式语义,并结合开源仓库源码(cmd/event-notification.go、internal/event/name.go、internal/event/arn.go、cmd/listen-notification-handlers.go)剖析事件匹配与发布的底层调用链,帮助你落地一套可复制、可验证的对象存储事件驱动架构。
一、事件模型:MinIO 支持哪些事件类型
MinIO 通过桶级事件通知机制,把桶内对象上发生的事件实时推送到外部系统。文档将事件类型划分为四类:
1. 对象事件(Object Events)
| 事件 | 事件 | 事件 |
|---|---|---|
s3:ObjectCreated:Put |
s3:ObjectCreated:CompleteMultipartUpload |
s3:ObjectAccessed:Head |
s3:ObjectCreated:Post |
s3:ObjectRemoved:Delete |
s3:ObjectRemoved:DeleteMarkerCreated |
s3:ObjectCreated:Copy |
s3:ObjectAccessed:Get |
|
s3:ObjectCreated:PutRetention |
s3:ObjectCreated:PutLegalHold |
|
s3:ObjectAccessed:GetRetention |
s3:ObjectAccessed:GetLegalHold |
2. 复制事件(Replication Events)
| 事件 |
|---|
s3:Replication:OperationFailedReplication |
s3:Replication:OperationCompletedReplication |
s3:Replication:OperationNotTracked |
s3:Replication:OperationMissedThreshold |
s3:Replication:OperationReplicatedAfterThreshold |
3. ILM 恢复事件(ILM Transition Events)
| 事件 |
|---|
s3:ObjectRestore:Post |
s3:ObjectRestore:Completed |
4. 全局事件(Global Events):仅通过 ListenNotification API 支持
| 事件 |
|---|
s3:BucketCreated |
s3:BucketRemoved |
从源码结构看,这些事件名在 internal/event/name.go 中定义为 Name 枚举(如 ObjectCreatedPut、ObjectRemovedDelete、ObjectReplicationFailed、ObjectRestorePost 等),且源码注释明确指出 s3:Replication:OperationCompletedReplication 属于 MinIO 扩展事件(非 AWS S3 原生语义)。Name 类型还提供 Expand() 方法,把复合类型展开为具体事件集:ObjectCreatedAll 展开为 CompleteMultipartUpload/Copy/Post/Put/PutRetention/PutLegalHold/PutTagging/DeleteTagging,ObjectRemovedAll 展开为 Delete/DeleteMarkerCreated/NoOP/DeleteAllVersions——这正是 mc event add 中 s3:ObjectCreated:* 这类通配写法能生效的底层机制。事件类型还采用位掩码设计(objectSingleTypesEnd 上限 64 个,用 uint64 位表示),保证 Match 匹配时的性能。
二、支持的通知目标总览
桶事件可以发布到以下十类目标,文档按目标逐一给出完整的三步配置法:
| 目标 | 目标 | 目标 |
|---|---|---|
| AMQP(RabbitMQ) | Redis | MySQL |
| MQTT | NATS | Apache Kafka |
| Elasticsearch | PostgreSQL | Webhooks |
| NSQ |
MinIO 客户端 mc 的 event 子命令用于设置与监听事件通知;各语言的 MinIO SDK 也提供 BucketNotification API。MinIO 发布的事件消息是符合 AWS S3 事件结构(eventVersion: 2.0)的 JSON 消息,字段结构见后文各目标的测试输出样例。
三、准备工作与十个通知子系统
前置条件:已安装并配置 MinIO Server 与 MinIO Client mc。用如下命令确认本部署支持的通知子系统:
$ mc admin config get myminio | grep notify
notify_webhook publish bucket notifications to webhook endpoints
notify_amqp publish bucket notifications to AMQP endpoints
notify_kafka publish bucket notifications to Kafka endpoints
notify_mqtt publish bucket notifications to MQTT endpoints
notify_nats publish bucket notifications to NATS endpoints
notify_nsq publish bucket notifications to NSQ endpoints
notify_mysql publish bucket notifications to MySQL databases
notify_postgres publish bucket notifications to Postgres databases
notify_elasticsearch publish bucket notifications to Elasticsearch endpoints
notify_redis publish bucket notifications to Redis datastores
配置约定(适用于下文中所有目标):
- 参数名末尾带
*表示必填; - 参数值末尾带
*表示该项是默认值; - 通过环境变量配置时,命名目标
:name用下划线后缀形式表达,例如MINIO_NOTIFY_WEBHOOK_ENABLE_<name>; - 每个目标可通过
notify_xxx:<name>追加任意多个端点实例,<name>(如1、myinstance)作为标识符。
配置修改流程统一为:mc admin config get myminio/ <target> 查看现状 → mc admin config set ... 更新 → 重启 MinIO Server 生效。启动成功无错误时,控制台会打印类似 SQS ARNs: arn:minio:sqs::1:amqp 的行,这是验证端点注册成功的标志。
3.1 ARN 格式与校验(源码佐证)
mc event add 引用的 ARN 遵循 arn:minio:sqs:<region>:<id>:<type> 格式。从 internal/event/arn.go 的 parseARN 函数可见:字符串必须以 arn:minio:sqs: 开头、按冒号切分后必须恰好 6 段且第 5、6 段非空,否则返回 ErrInvalidARN;其中 <id> 对应配置名(如 1、myinstance),<type> 对应子系统类型(如 amqp、kafka)。这也是为什么自定义 notify_mysql:myinstance 后,其 ARN 写作 arn:minio:sqs::myinstance:mysql。
3.2 事件发布调用链(源码佐证)
cmd/event-notification.go 是通知系统的核心:
EventNotifier持有targetList(全部外部目标)与bucketRulesMap(每个桶的事件规则表);- 对象操作完成时调用
sendEvent(args eventArgs):先跳过复制源端请求(MinIOSourceReplicationRequest,避免副本侧重复触发)、剥离加密敏感元数据;若存在 HTTP 长连接订阅者(ListenNotification)则通过globalHTTPListen.Publish推送; - 随后
EventNotifier.Send用bucketRulesMap[bucket].Match(eventName, objectName)按桶规则+前缀/后缀过滤器算出目标集合targetIDSet,无匹配则直接返回; - 事件体由
ToEvent构造:EventVersion: "2.0"、EventSource: "minio:s3"、Sequencer取对象修改时间的纳秒十六进制值,响应头包含x-amz-request-id、x-minio-origin-endpoint、x-minio-deployment-id等 MinIO 自定义字段;删除类事件(Delete/DeleteMarkerCreated/NoOP)不携带ETag/Size/ContentType。 - 是否同步发送由 API 配置
MINIO_API_SYNC_EVENTS控制(见Send中isSyncEventsEnabled())。
3.3 ListenNotification API
除了外部目标推送,MinIO 还提供 S3 兼容的 ListenNotification API(实现在 cmd/listen-notification-handlers.go):客户端发起长连接请求,服务端按 prefix/suffix 过滤条件把匹配事件以 chunk 形式持续写回。权限上,指定桶名走 ListenBucketNotification 动作,桶名为空则要求 ListenNotification 动作(全局),文档中列出的 s3:BucketCreated/s3:BucketRemoved 全局事件仅在此通道可见。SDK 中对应的 ListenBucketNotification 方法适合编写轻量消费端,无需部署 Kafka/RabbitMQ 等中间件。
3.4 持久化事件存储(queue_dir / queue_limit)
所有消息队列/存储类目标共享一套持久事件存储机制:目标端离线时,事件备份到本地磁盘目录,目标端恢复后自动回放。两个参数在所有目标中语义一致:
| 参数 | 类型 | 说明 |
|---|---|---|
queue_dir |
path | 未投递事件的暂存目录,例如 /home/events |
queue_limit |
number | 暂存事件数量上限,默认 100000 |
四、AMQP(RabbitMQ)
安装 RabbitMQ 后即可接入。
4.1 添加 AMQP 端点
配置位于 notify_amqp 键下。参数表:
KEY:
notify_amqp[:name] publish bucket notifications to AMQP endpoints
ARGS:
url (url) AMQP server endpoint e.g. `amqp://myuser:mypassword@localhost:5672`
exchange (string) name of the AMQP exchange
exchange_type (string) AMQP exchange type
routing_key (string) routing key for publishing
mandatory (on|off) quietly ignore undelivered messages when set to 'off', default is 'on'
durable (on|off) persist queue across broker restarts when set to 'on', default is 'off'
no_wait (on|off) non-blocking message delivery when set to 'on', default is 'off'
internal (on|off) set to 'on' for exchange to be not used directly by publishers, but only when bound to other exchanges
auto_deleted (on|off) auto delete queue when set to 'on', when there are no consumers
delivery_mode (number) set to '1' for non-persistent or '2' for persistent queue
queue_dir (path) staging dir for undelivered messages e.g. '/home/events'
queue_limit (number) maximum limit for undelivered messages, defaults to '100000'
comment (sentence) optionally add a comment to this setting
(url 为必填项。环境变量形式即把上述参数映射为 MINIO_NOTIFY_AMQP_ENABLE*、MINIO_NOTIFY_AMQP_URL*、MINIO_NOTIFY_AMQP_EXCHANGE 等,命名规则见第三节约定。)
$ mc admin config get myminio/ notify_amqp
notify_amqp:1 delivery_mode="0" exchange_type="" no_wait="off" queue_dir="" queue_limit="0" url="" auto_deleted="off" durable="off" exchange="" internal="off" mandatory="off" routing_key=""
示例配置(fanout 交换机 bucketevents,routing key bucketlogs):
mc admin config set myminio/ notify_amqp:1 exchange="bucketevents" exchange_type="fanout" mandatory="off" no_wait="off" url="amqp://myuser:mypassword@localhost:5672" auto_deleted="off" delivery_mode="0" durable="off" internal="off" routing_key="bucketlogs"
要点:MinIO 支持 RabbitMQ 的全部交换机类型;MinIO 还随通知附带 minio-bucket、minio-event 两个消息头,headers 类型交换机可据此把通知路由到正确队列。
4.2 启用桶通知
对 images 桶按后缀过滤 JPEG 上传/删除事件,ARN 为 arn:minio:sqs::1:amqp:
mc mb myminio/images
mc event add myminio/images arn:minio:sqs::1:amqp --suffix .jpg
mc event list myminio/images
arn:minio:sqs::1:amqp s3:ObjectCreated:*,s3:ObjectRemoved:* Filter: suffix=".jpg"
4.3 用 Python 消费者验证
用 Pika 客户端等待交换机 bucketevents 并打印事件:
#!/usr/bin/env python
import pika
connection = pika.BlockingConnection(pika.ConnectionParameters(
host='localhost'))
channel = connection.channel()
channel.exchange_declare(exchange='bucketevents',
exchange_type='fanout')
result = channel.queue_declare(exclusive=False)
queue_name = result.method.queue
channel.queue_bind(exchange='bucketevents',
queue=queue_name)
print(' [*] Waiting for logs. To exit press CTRL+C')
def callback(ch, method, properties, body):
print(" [x] %r" % body)
channel.basic_consume(callback,
queue=queue_name,
no_ack=False)
channel.start_consuming()
另开终端上传 JPEG:mc cp myphoto.jpg myminio/images。上传完成后即收到形如以下的事件:
{"Records":[{"eventVersion":"2.0","eventSource":"aws:s3","awsRegion":"","eventTime":"2016-09-08T22:34:38.226Z","eventName":"s3:ObjectCreated:Put","userIdentity":{"principalId":"minio"},"requestParameters":{"sourceIPAddress":"10.1.10.150:44576"},"responseElements":{},"s3":{"s3SchemaVersion":"1.0","configurationId":"Config","bucket":{"name":"images","ownerIdentity":{"principalId":"minio"},"arn":"arn:aws:s3:::images"},"object":{"key":"myphoto.jpg","size":200436,"sequencer":"147279EAF9F40933"}}}],"level":"info","msg":"","time":"2016-09-08T15:34:38-07:00"}
五、MQTT
安装 MQTT Broker(如 Mosquitto)后即可接入。
5.1 添加 MQTT 端点
KEY:
notify_mqtt[:name] publish bucket notifications to MQTT endpoints
ARGS:
broker (uri) MQTT server endpoint e.g. `tcp://localhost:1883`
topic (string) name of the MQTT topic to publish
username (string) MQTT username
password (string) MQTT password
qos (number) set the quality of service priority, defaults to '0'
keep_alive_interval (duration) keep-alive interval for MQTT connections in s,m,h,d
reconnect_interval (duration) reconnect interval for MQTT connections in s,m,h,d
queue_dir (path) staging dir for undelivered messages e.g. '/home/events'
queue_limit (number) maximum limit for undelivered messages, defaults to '100000'
comment (sentence) optionally add a comment to this setting
(broker、topic 必填。环境变量前缀 MINIO_NOTIFY_MQTT_*。)
MinIO 支持 MQTT 3.1 / 3.1.1 协议的任意 Broker,broker URL 分别以 tcp://、tls://、ws:// 作为 scheme 表示 TCP、TLS 或 WebSocket 连接。
$ mc admin config get myminio/ notify_mqtt
notify_mqtt:1 broker="" password="" queue_dir="" queue_limit="0" reconnect_interval="0s" keep_alive_interval="0s" qos="0" topic="" username=""
mc admin config set myminio notify_mqtt:1 broker="tcp://localhost:1883" password="" queue_dir="" queue_limit="0" reconnect_interval="0s" keep_alive_interval="0s" qos="1" topic="minio" username=""
5.2 启用桶通知
ARN 为 arn:minio:sqs::1:mqtt:
mc mb myminio/images
mc event add myminio/images arn:minio:sqs::1:mqtt --suffix .jpg
mc event list myminio/images
5.3 用 paho-mqtt 验证
#!/usr/bin/env python3
from __future__ import print_function
import paho.mqtt.client as mqtt
# This is the Subscriber
def on_connect(client, userdata, flags, rc):
print("Connected with result code "+str(rc))
# qos level is set to 1
client.subscribe("minio", 1)
def on_message(client, userdata, msg):
print(msg.payload)
# client_id is a randomly generated unique ID for the mqtt broker to identify the connection.
client = mqtt.Client(client_id="myclientid", clean_session=False)
client.on_connect = on_connect
client.on_message = on_message
client.connect("localhost",1883,60)
client.loop_forever()
运行 python mqtt.py 后执行 mc cp myphoto.jpg myminio/images,即收到与 AMQP 样例同构的 s3:ObjectCreated:Put JSON 事件(含 Records[0].s3.object.key、size、sequencer 等字段)。
六、Elasticsearch
安装 Elasticsearch 后接入。此目标支持 namespace 与 access 两种格式:
- namespace:桶内对象与索引文档双向同步——以
桶名/对象名作为文档 ID;对象覆写时更新文档,对象删除时删除文档。索引内容始终反映桶的当前状态; - access:把事件追加为文档(ID 由 Elasticsearch 随机生成,文档时间戳取事件时间),不删除也不修改任何已有文档,形成操作日志流。
以下步骤以 namespace 格式为例。
6.1 版本要求
MinIO 要求 Elasticsearch 5.x 系列。
6.2 添加 Elasticsearch 端点
KEY:
notify_elasticsearch[:name] publish bucket notifications to Elasticsearch endpoints
ARGS:
url (url) Elasticsearch server's address, with optional authentication info
index (string) Elasticsearch index to store/update events, index is auto-created
format (namespace*|access) 'namespace' reflects current bucket/object list and 'access' reflects a journal of object operations, defaults to 'namespace'
queue_dir (path) staging dir for undelivered messages e.g. '/home/events'
queue_limit (number) maximum limit for undelivered messages, defaults to '100000'
username (string) username for Elasticsearch basic-auth
password (string) password for Elasticsearch basic-auth
comment (sentence) optionally add a comment to this setting
(url、index、format 必填,环境变量前缀 MINIO_NOTIFY_ELASTICSEARCH_*。)
url 示例:http://localhost:9200,或带 basic-auth 的 http://elastic:MagicWord@127.0.0.1:9200——认证信息直接编码在 URL 中。
$ mc admin config get myminio/ notify_elasticsearch
notify_elasticsearch:1 queue_limit="0" url="" format="namespace" index="" queue_dir=""
mc admin config set myminio notify_elasticsearch:1 queue_limit="0" url="http://127.0.0.1:9200" format="namespace" index="minio_events" queue_dir="" username="" password=""
6.3 启用桶通知
ARN 为 arn:minio:sqs::1:elasticsearch:
mc mb myminio/images
mc event add myminio/images arn:minio:sqs::1:elasticsearch --suffix .jpg
mc event list myminio/images
arn:minio:sqs::1:elasticsearch s3:ObjectCreated:*,s3:ObjectRemoved:* Filter: suffix=".jpg"
此后 minio_events 索引中的文档集合即与 images 桶中的 .jpg 对象一一对应。
6.4 用 curl 验证
mc cp myphoto.jpg myminio/images
curl "http://localhost:9200/minio_events/_search?pretty=true"
返回结果关键字段:
{
"hits" : {
"total" : 1,
"hits" : [
{
"_index" : "minio_events",
"_type" : "event",
"_id" : "images/myphoto.jpg",
"_source" : {
"Records" : [
{
"eventVersion" : "2.0",
"eventSource" : "minio:s3",
"eventName" : "s3:ObjectCreated:Put",
"s3" : {
"bucket" : { "name" : "images", "arn" : "arn:aws:s3:::images" },
"object" : {
"key" : "myphoto.jpg",
"size" : 6474,
"eTag" : "a3410f4f8788b510d6f19c5067e60a90",
"sequencer" : "14B09A09703FC47B"
}
},
"source" : {
"host" : "127.0.0.1",
"userAgent" : "MinIO (linux; amd64) minio-go/2.0.3 mc/2017-02-15T17:57:25Z"
}
}
]
}
}
]
}
}
文档 ID images/myphoto.jpg 正是 namespace 格式的 桶名/对象名;若使用 access 格式,文档 ID 会由 Elasticsearch 自动生成。
七、Redis
安装 Redis 后接入(示例中密码设为 yoursecret)。此目标同样支持两种格式:
- namespace:桶内对象与一个 hash 同步。键为
bucketName/objectName,值为创建/覆写该对象的操作的 JSON 事件;对象更新或删除时对应键被更新或删除; - access:用
RPUSH把事件追加到 list,每个元素是[时间戳字符串, 事件 JSON]两元素列表;追加后 MinIO 不再修改或删除。
7.1 添加 Redis 端点
KEY:
notify_redis[:name] publish bucket notifications to Redis datastores
ARGS:
address (address) Redis server's address. For example: `localhost:6379`
key (string) Redis key to store/update events, key is auto-created
format (namespace*|access) 'namespace' reflects current bucket/object list and 'access' reflects a journal of object operations, defaults to 'namespace'
password (string) Redis server password
queue_dir (path) staging dir for undelivered messages e.g. '/home/events'
queue_limit (number) maximum limit for undelivered messages, defaults to '100000'
comment (sentence) optionally add a comment to this setting
(address、key、format 必填,环境变量前缀 MINIO_NOTIFY_REDIS_*。)
$ mc admin config get myminio/ notify_redis
notify_redis:1 address="" format="namespace" key="" password="" queue_dir="" queue_limit="0"
mc admin config set myminio/ notify_redis:1 address="127.0.0.1:6379" format="namespace" key="bucketevents" password="yoursecret" queue_dir="" queue_limit="0"
7.2 启用桶通知
ARN 为 arn:minio:sqs::1:redis:
mc mb myminio/images
mc event add myminio/images arn:minio:sqs::1:redis --suffix .jpg
mc event list myminio/images
arn:minio:sqs::1:redis s3:ObjectCreated:*,s3:ObjectRemoved:* Filter: suffix=".jpg"
7.3 用 redis-cli monitor 验证
redis-cli -a yoursecret
127.0.0.1:6379> monitor
OK
另开终端 mc cp myphoto.jpg myminio/images,monitor 输出将显示 MinIO 在 Redis 上执行的操作:
1490686879.650649 [0 172.17.0.1:44710] "PING"
1490686879.651061 [0 172.17.0.1:44710] "HSET" "minio_events" "images/myphoto.jpg" "{\"Records\":[{\"eventVersion\":\"2.0\",\"eventSource\":\"minio:s3\",...\"object\":{\"key\":\"myphoto.jpg\",\"size\":2586,\"eTag\":\"5d284463f9da279f060f0ea4d11af098\",...}}]}"
namespace 格式对应 HSET 写 hash;若使用 access 格式,则 minio_events 是一个 list,MinIO 执行 RPUSH 追加,消费端理想地用 BLPOP 从左端取走条目。
八、NATS
安装 NATS 后接入。除普通 NATS 外还支持 NATS Streaming 模式,额外提供至少一次投递(At-least-once-delivery)与发布端限流(Publisher rate limiting)。
8.1 添加 NATS 端点(参数最全的目标之一)
KEY:
notify_nats[:name] publish bucket notifications to NATS endpoints
ARGS:
address (address) NATS server address e.g. '0.0.0.0:4222'
subject (string) NATS subscription subject
username (string) NATS username
password (string) NATS password
token (string) NATS token
tls (on|off) set to 'on' to enable TLS
tls_skip_verify (on|off) trust server TLS without verification, defaults to "on" (verify)
ping_interval (duration) client ping commands interval in s,m,h,d. Disabled by default
streaming (on|off) set to 'on', to use streaming NATS server
streaming_async (on|off) set to 'on', to enable asynchronous publish
streaming_max_pub_acks_in_flight (number) number of messages to publish without waiting for ACKs
streaming_cluster_id (string) unique ID for NATS streaming cluster
cert_authority (string) path to certificate chain of the target NATS server
client_cert (string) client cert for NATS mTLS auth
client_key (string) client cert key for NATS mTLS auth
queue_dir (path) staging dir for undelivered messages e.g. '/home/events'
queue_limit (number) maximum limit for undelivered messages, defaults to '100000'
comment (sentence) optionally add a comment to this setting
(address、subject 必填;环境变量前缀 MINIO_NOTIFY_NATS_*,含 MINIO_NOTIFY_NATS_STREAMING、MINIO_NOTIFY_NATS_STREAMING_CLUSTER_ID、MINIO_NOTIFY_NATS_CERT_AUTHORITY 等。)
$ mc admin config get myminio/ notify_nats
notify_nats:1 password="yoursecret" streaming_max_pub_acks_in_flight="10" subject="" address="0.0.0.0:4222" token="" username="yourusername" ping_interval="0" queue_limit="0" tls="off" tls_skip_verify="off" streaming_async="on" queue_dir="" streaming_cluster_id="test-cluster" streaming_enable="on"
mc admin config set myminio notify_nats:1 password="yoursecret" streaming_max_pub_acks_in_flight="10" subject="" address="0.0.0.0:4222" token="" username="yourusername" ping_interval="0" queue_limit="0" tls="off" streaming_async="on" queue_dir="" streaming_cluster_id="test-cluster" streaming_enable="on"
配置 NATS Streaming 时,关键为 streaming="on"、streaming_cluster_id(集群 ID)与 streaming_max_pub_acks_in_flight(不等 ACK 的并发发布数,即发布端限流值),并配合 cert_authority/client_cert/client_key 做 mTLS。
8.2 启用桶通知
ARN 为 arn:minio:sqs::1:nats:
mc mb myminio/images
mc event add myminio/images arn:minio:sqs::1:nats --suffix .jpg
mc event list myminio/images
arn:minio:sqs::1:nats s3:ObjectCreated:*,s3:ObjectRemoved:* Filter: suffix=".jpg"
8.3 订阅验证
普通 NATS 的 Go 订阅示例(订阅 subject bucketevents):
package main
import (
"log"
"runtime"
"github.com/nats-io/nats.go"
)
func main() {
// Create server connection
natsConnection, _ := nats.Connect("nats://yourusername:yoursecret@localhost:4222")
log.Println("Connected")
// Subscribe to subject
log.Printf("Subscribing to subject 'bucketevents'\n")
natsConnection.Subscribe("bucketevents", func(msg *nats.Msg) {
// Handle the message
log.Printf("Received message '%s\n", string(msg.Data)+"'")
})
// Keep the connection alive
runtime.Goexit()
}
上传 mc cp myphoto.jpg myminio/images 后,控制台输出形如:
Received message '{"EventType":"s3:ObjectCreated:Put","Key":"images/myphoto.jpg","Records":[{"eventVersion":"2.0","eventSource":"aws:s3",...,"eventName":"s3:ObjectCreated:Put",...,"object":{"key":"myphoto.jpg","size":56060,"eTag":"1d97bf45ecb37f7a7b699418070df08f","sequencer":"147CCD1AE054BFD0"}}],...}'
NATS Streaming 版本则用 stan.go:stan.Connect("test-cluster", "test-client", stan.NatsURL(...), stan.SetConnectionLostHandler(...)) 建立连接并订阅 bucketevents,断线时在 SetConnectionLostHandler 回调中重连。收到的事件 JSON 结构与上面一致,且包含 contentType、versionId、userDefined 等更完整的对象字段。
九、PostgreSQL
升级注意:在 RELEASE.2020-04-10T03-34-42Z 之前,PostgreSQL 通知还支持
host/port/username/password/database五个独立选项;这些选项已废弃。升级到该版本之后的发布前,必须迁移为仅使用connection_string:mc admin config set myminio/ notify_postgres[:name] connection_string="host=hostname port=5432 username=psqluser password=psqlpass database=bucketevents"不做此迁移会导致 PostgreSQL 通知目标失效,升级/重启时控制台会打印报错。
示例环境:postgres 用户密码为 password,已创建数据库 minio_events。
此目标支持两种格式:
- namespace:桶内对象与表行同步。表含
key/value两列:key 为存在的对象的桶名+对象名,value 为创建/覆写操作的 JSON 事件;对象更新/删除时对应行更新/删除; - access:把事件追加为表行,两列
event_time(事件在 MinIO 服务端发生的时间)与event_data(事件 JSON),不删除也不修改行。
以下以 namespace 格式为例。
9.1 版本要求
MinIO 要求 PostgreSQL 9.5+,依赖 9.5 引入的 INSERT ON CONFLICT(UPSERT)特性与 9.4 引入的 JSONB 类型。
9.2 添加 PostgreSQL 端点
KEY:
notify_postgres[:name] publish bucket notifications to Postgres databases
ARGS:
connection_string (string) Postgres server connection-string e.g. "host=localhost port=5432 dbname=minio_events user=postgres password=password sslmode=disable"
table (string) DB table name to store/update events, table is auto-created
format (namespace*|access) 'namespace' reflects current bucket/object list and 'access' reflects a journal of object operations, defaults to 'namespace'
queue_dir (path) staging dir for undelivered messages e.g. '/home/events'
queue_limit (number) maximum limit for undelivered messages, defaults to '100000'
max_open_connections (number) maximum number of open connections to the database, defaults to '2'
comment (sentence) optionally add a comment to this setting
(前三项必填;环境变量前缀 MINIO_NOTIFY_POSTGRES_*。注意:max_open_connections 设为 0 表示不限制连接数——官方明确不推荐,namespace 格式下的递归删除可能出现不一致行为。)
$ mc admin config get myminio notify_postgres
notify_postgres:1 queue_dir="" connection_string="" queue_limit="0" table="" format="namespace"
mc admin config set myminio notify_postgres:1 connection_string="host=localhost port=5432 dbname=minio_events user=postgres password=password sslmode=disable" table="bucketevents" format="namespace"
示例中禁用了 SSL 仅为演示,生产环境不建议。
9.3 启用桶通知
ARN 为 arn:minio:sqs::1:postgresql:
# Create bucket named `images` in myminio
mc mb myminio/images
# Add notification configuration on the `images` bucket using the ARN. The --suffix argument filters events.
mc event add myminio/images arn:minio:sqs::1:postgresql --suffix .jpg
# Print out the notification configuration on the `images` bucket.
mc event list myminio/images
arn:minio:sqs::1:postgresql s3:ObjectCreated:*,s3:ObjectRemoved:* Filter: suffix=".jpg"
9.4 用 psql 验证
mc cp myphoto.jpg myminio/images
psql -h 127.0.0.1 -U postgres -d minio_events
minio_events=# select * from bucketevents;
key | value
--------------------+-----------------------------------------------------------------
images/myphoto.jpg | {"Records": [{"s3": {"bucket": {"arn": "arn:aws:s3:::images", "name": "images", ...}, "object": {"key": "myphoto.jpg", "eTag": "1d97bf45ecb37f7a7b699418070df08f", "size": 56060, "sequencer": "147CE57C70B31931"}, ...}, "eventName": "s3:ObjectCreated:Put", "eventTime": "2016-10-12T21:18:20Z", "eventSource": "aws:s3", ...}]}
(1 row)
十、MySQL
升级注意:与 PostgreSQL 同理,RELEASE.2020-04-10T03-34-42Z 之前的
host/port/username/password/database独立选项已废弃,升级后必须迁移为dsn_string:mc admin config set myminio/ notify_mysql[:name] dsn_string="mysqluser:mysqlpass@tcp(localhost:3306)/bucketevents"不迁移则 MySQL 通知目标失效,控制台升级/重启时会有报错。
示例环境:root 密码为 password,已创建数据库 miniodb。
此目标支持两种格式:
- namespace:桶内对象与表行同步,两列
key_name(桶名+对象名)与value(JSON 事件);对象更新/删除时行相应更新/删除; - access:追加
event_time+event_data行,不删不改。
以下以 namespace 格式为例。
10.1 版本要求
MinIO 要求 MySQL 5.7.8+(依赖 5.7.8 引入的 JSON 类型)。文档中该组合在 MySQL 5.7.17 上做过验证。
10.2 添加 MySQL 端点
KEY:
notify_mysql[:name] publish bucket notifications to MySQL databases. When multiple MySQL server endpoints are needed, a user specified "name" can be added for each configuration, (e.g."notify_mysql:myinstance").
ARGS:
dsn_string (string) MySQL data-source-name connection string e.g. "<user>:<password>@tcp(<host>:<port>)/<database>"
table (string) DB table name to store/update events, table is auto-created
format (namespace*|access) 'namespace' reflects current bucket/object list and 'access' reflects a journal of object operations, defaults to 'namespace'
queue_dir (path) staging dir for undelivered messages e.g. '/home/events'
queue_limit (number) maximum limit for undelivered messages, defaults to '100000'
max_open_connections (number) maximum number of open connections to the database, defaults to '2'
comment (sentence) optionally add a comment to this setting
(前三项必填;环境变量前缀 MINIO_NOTIFY_MYSQL_*,含 MINIO_NOTIFY_MYSQL_DSN_STRING*、MINIO_NOTIFY_MYSQL_MAX_OPEN_CONNECTIONS 等。max_open_connections=0 不推荐,理由同 PostgreSQL。)
$ mc admin config get myminio/ notify_mysql
notify_mysql:myinstance enable=off format=namespace host= port= username= password= database= dsn_string= table= queue_dir= queue_limit=0
mc admin config set myminio notify_mysql:myinstance table="minio_images" dsn_string="root:xxxx@tcp(172.17.0.1:3306)/miniodb"
10.3 启用桶通知
自定义名 myinstance 的 ARN 为 arn:minio:sqs::myinstance:mysql:
# Create bucket named `images` in myminio
mc mb myminio/images
# Add notification configuration on the `images` bucket using the MySQL ARN. The --suffix argument filters events.
mc event add myminio/images arn:minio:sqs::myinstance:mysql --suffix .jpg
# Print out the notification configuration on the `images` bucket.
mc event list myminio/images
arn:minio:sqs::myinstance:mysql s3:ObjectCreated:*,s3:ObjectRemoved:*,s3:ObjectAccessed:* Filter: suffix=".jpg"
10.4 用 mysql 客户端验证
mc cp myphoto.jpg myminio/images
mysql -h 172.17.0.1 -P 3306 -u root -p miniodb
mysql> select * from minio_images;
+--------------------+-----------------------------------------------------------------------------------+
| key_name | value |
+--------------------+-----------------------------------------------------------------------------------+
| images/myphoto.jpg | {"Records": [{"s3": {"bucket": {"arn": "arn:aws:s3:::images", "name": "images", ...}, "object": {"key": "myphoto.jpg", "eTag": "467886be95c8ecfd71a2900e3f461b4f", "size": 26, "sequencer": "14AC59476F809FD3"}, ...}, "eventName": "s3:ObjectCreated:Put", ...}]} |
+--------------------+-----------------------------------------------------------------------------------+
1 row in set (0.01 sec)
十一、Apache Kafka
11.1 版本要求
MinIO 依赖 sarama 库(Shopify/sarama),Kafka 版本兼容性与该库一致,文档以 0.9/0.10+ 集群为基准。
11.2 添加 Kafka 端点
KEY:
notify_kafka[:name] publish bucket notifications to Kafka endpoints
ARGS:
brokers (csv) comma separated list of Kafka broker addresses
topic (string) Kafka topic used for bucket notifications
sasl_username (string) username for SASL/PLAIN or SASL/SCRAM authentication
sasl_password (string) password for SASL/PLAIN or SASL/SCRAM authentication
sasl_mechanism (string) sasl authentication mechanism, default 'PLAIN'
tls_client_auth (string) clientAuth determines the Kafka server's policy for TLS client auth
sasl (on|off) set to 'on' to enable SASL authentication
tls (on|off) set to 'on' to enable TLS
tls_skip_verify (on|off) trust server TLS without verification, defaults to "on" (verify)
client_tls_cert (path) path to client certificate for mTLS auth
client_tls_key (path) path to client key for mTLS auth
queue_dir (path) staging dir for undelivered messages e.g. '/home/events'
queue_limit (number) maximum limit for undelivered messages, defaults to '100000'
version (string) specify the version of the Kafka cluster e.g '2.2.0'
comment (sentence) optionally add a comment to this setting
(brokers 必填;环境变量前缀 MINIO_NOTIFY_KAFKA_*,此外还有两个仅环境变量形式的生产者压缩选项:MINIO_NOTIFY_KAFKA_PRODUCER_COMPRESSION_CODEC(none|snappy|gzip|lz4|zstd)与 MINIO_NOTIFY_KAFKA_PRODUCER_COMPRESSION_LEVEL(默认 0)。)
$ mc admin config get myminio/ notify_kafka
notify_kafka:1 tls_skip_verify="off" queue_dir="" queue_limit="0" sasl="off" sasl_password="" sasl_username="" tls_client_auth="0" tls="off" brokers="" topic="" client_tls_cert="" client_tls_key="" version=""
mc admin config set myminio notify_kafka:1 tls_skip_verify="off" queue_dir="" queue_limit="0" sasl="off" sasl_password="" sasl_username="" tls_client_auth="0" tls="off" client_tls_cert="" client_tls_key="" brokers="localhost:9092,localhost:9093" topic="bucketevents" version=""
brokers 支持逗号分隔多 Broker;SASL 机制可选 plain|sha256|sha512;mTLS 走 client_tls_cert/client_tls_key。
11.3 启用桶通知
ARN 为 arn:minio:sqs::1:kafka:
mc mb myminio/images
mc event add myminio/images arn:minio:sqs::1:kafka --suffix .jpg
mc event list myminio/images
arn:minio:sqs::1:kafka s3:ObjectCreated:*,s3:ObjectRemoved:* Filter: suffix=".jpg"
11.4 用 kafkacat 验证
kafkacat -C -b localhost:9092 -t bucketevents
另开终端 mc cp myphoto.jpg myminio/images,kafkacat 打印出完整事件,示例关键字段:
{
"EventName": "s3:ObjectCreated:Put",
"Key": "images/myphoto.jpg",
"Records": [
{
"eventVersion": "2.0",
"eventSource": "minio:s3",
"eventTime": "2019-09-10T17:41:54Z",
"eventName": "s3:ObjectCreated:Put",
"userIdentity": { "principalId": "AKIAIOSFODNN7EXAMPLE" },
"requestParameters": { "accessKey": "AKIAIOSFODNN7EXAMPLE", "sourceIPAddress": "192.168.56.192" },
"responseElements": {
"x-amz-request-id": "15C3249451E12784",
"x-minio-deployment-id": "751a8ba6-acb2-42f6-a297-4cdf1cf1fa4f",
"x-minio-origin-endpoint": "http://192.168.97.83:9000"
},
"s3": {
"s3SchemaVersion": "1.0",
"bucket": { "name": "images", "arn": "arn:aws:s3:::images" },
"object": {
"key": "myphoto.jpg",
"size": 6474,
"eTag": "430f89010c77aa34fc8760696da62d08-1",
"contentType": "image/jpeg",
"userMetadata": { "content-type": "image/jpeg" },
"versionId": "1",
"sequencer": "15C32494527B46C5"
}
},
"source": { "host": "192.168.56.192", "userAgent": "Mozilla/5.0 ..." }
}
]
}
十二、Webhooks
Webhook 是"事件发生时主动推送"而非轮询的集成方式,也是把对象存储接入缩略图生成、索引构建等自研服务的最低成本通道。
12.1 添加 Webhook 端点
KEY:
notify_webhook[:name] publish bucket notifications to webhook endpoints
ARGS:
endpoint (url) webhook server endpoint e.g. http://localhost:8080/minio/events
auth_token (string) opaque string or JWT authorization token
queue_dir (path) staging dir for undelivered messages e.g. '/home/events'
queue_limit (number) maximum limit for undelivered messages, defaults to '100000'
client_cert (string) client cert for Webhook mTLS auth
client_key (string) client cert key for Webhook mTLS auth
comment (sentence) optionally add a comment to this setting
(endpoint 必填;环境变量前缀 MINIO_NOTIFY_WEBHOOK_*。)
$ mc admin config get myminio/ notify_webhook
notify_webhook:1 endpoint="" auth_token="" queue_limit="0" queue_dir="" client_cert="" client_key=""
mc admin config set myminio notify_webhook:1 queue_limit="0" endpoint="http://localhost:3000" queue_dir=""
注意:重启 MinIO 时该 endpoint 必须在线且可达。auth_token 可承载不透明字符串或 JWT;client_cert/client_key 用于对目标端做 mTLS。
12.2 启用桶通知
mc mb myminio/images
mc mb myminio/images-thumbnail
mc event add myminio/images arn:minio:sqs::1:webhook --event put --suffix .jpg
mc event list myminio/images
arn:minio:sqs::1:webhook s3:ObjectCreated:* Filter: suffix=".jpg"
12.3 以 Thumbnailer 做端到端验证
用 minio/thumbnailer 项目监听通知:新 JPEG 上传(PUT)触发后,Thumbnailer 把缩略图回传到 MinIO。步骤:
git clone https://github.com/minio/thumbnailer/
npm install
在 config/webhook.json 中填写 MinIO 服务器配置,然后启动:
NODE_ENV=webhook node thumbnail-webhook.js
Thumbnailer 运行于 http://localhost:3000/(与 12.1 中 endpoint 对应)。上传一张 JPEG:
mc cp ~/images.jpg myminio/images
.../images.jpg: 8.31 KB / 8.31 KB 100.00%
稍等片刻后查看缩略图桶:
mc ls myminio/images-thumbnail
[2017-02-08 11:39:40 IST] 992B images-thumbnail.jpg
十三、NSQ
安装 nsq 守护进程,或用容器方式启动:
podman run --rm -p 4150-4151:4150-4151 nsqio/nsq /nsqd
13.1 添加 NSQ 端点
KEY:
notify_nsq[:name] publish bucket notifications to NSQ endpoints
ARGS:
nsqd_address (address) NSQ server address e.g. '127.0.0.1:4150'
topic (string) NSQ topic
tls (on|off) set to 'on' to enable TLS
tls_skip_verify (on|off) trust server TLS without verification, defaults to "on" (verify)
queue_dir (path) staging dir for undelivered messages e.g. '/home/events'
queue_limit (number) maximum limit for undelivered messages, defaults to '100000'
comment (sentence) optionally add a comment to this setting
(nsqd_address、topic 必填;环境变量前缀 MINIO_NOTIFY_NSQ_*。)
$ mc admin config get myminio/ notify_nsq
notify_nsq:1 nsqd_address="" queue_dir="" queue_limit="0" tls="off" tls_skip_verify="off" topic=""
mc admin config set myminio notify_nsq:1 nsqd_address="127.0.0.1:4150" queue_dir="" queue_limit="0" tls="off" tls_skip_verify="on" topic="minio"
13.2 启用桶通知
ARN 为 arn:minio:sqs::1:nsq:
mc mb myminio/images
mc event add myminio/images arn:minio:sqs::1:nsq --suffix .jpg
mc event list myminio/images
arn:minio:sqs::1:nsq s3:ObjectCreated:*,s3:ObjectRemoved:* Filter: suffix=".jpg"
13.3 用 nsq_tail 验证
从 NSQ 发行版下载 nsq_tail:
./nsq_tail -nsqd-tcp-address 127.0.0.1:4150 -topic minio
另开终端 mc cp gopher.jpg myminio/images,nsq_tail 输出:
{"EventName":"s3:ObjectCreated:Put","Key":"images/gopher.jpg","Records":[{"eventVersion":"2.0","eventSource":"minio:s3","awsRegion":"","eventTime":"2018-10-31T09:31:11Z","eventName":"s3:ObjectCreated:Put","userIdentity":{"principalId":"21EJ9HYV110O8NVX2VMS"},"requestParameters":{"sourceIPAddress":"10.1.1.1"},"responseElements":{"x-amz-request-id":"1562A792DAA53426","x-minio-origin-endpoint":"http://10.0.3.1:9000"},"s3":{"s3SchemaVersion":"1.0","configurationId":"Config","bucket":{"name":"images","ownerIdentity":{"principalId":"21EJ9HYV110O8NVX2VMS"},"arn":"arn:aws:s3:::images"},"object":{"key":"gopher.jpg","size":162023,"eTag":"5337769ffa594e742408ad3f30713cd7","contentType":"image/jpeg","userMetadata":{"content-type":"image/jpeg"},"versionId":"1","sequencer":"1562A792DAA53426"}},"source":{"host":"","port":"","userAgent":"MinIO (linux; amd64) minio-go/v6.0.8 mc/DEVELOPMENT.GOGET"}}]}
十四、实践要点与源码索引
- 过滤粒度:
mc event add支持--event(指定事件,如put)、--prefix/--suffix(前缀/后缀过滤)。过滤器校验逻辑在 internal/event/rules.go;每个桶的规则集以RulesMap(规则 → 目标集合)形式缓存在EventNotifier.bucketRulesMap中,热路径Send只需一次加读锁匹配。 - 多端点:任何目标都可通过
notify_xxx:<name>注册多实例,事件会被同时推送到 ARN 引用的所有端点;internal/event/targetid.go 定义了TargetID及其到 ARN 字符串的互转。 - 可靠性分层:队列类目标(AMQP/MQTT/NATS/Kafka/NSQ/Webhook)与存储类目标(Redis/ES/PostgreSQL/MySQL)都可用
queue_dir/queue_limit启用磁盘暂存,目标端恢复后自动重放;queue_limit默认 100000 条。 - 语义差异:消息队列类目标是"事件流"(append-only);而 Redis/ES/PostgreSQL/MySQL 的
namespace格式是"状态镜像"(upsert/delete 与桶内对象生命周期一致),access格式则退化为"事件日志"。选型时按消费端是否需要全量当前状态来区分。 - 可观测性:事件体中的
x-amz-request-id、x-minio-deployment-id、x-minio-origin-endpoint可与访问日志、复制指标交叉关联,便于审计链路排查(字段构造见 cmd/event-notification.go 的ToEvent)。 - 相关文档与源码索引:
| 主题 | 路径 |
|---|---|
| 本文档原文 | docs/bucket/notifications/README.md |
| 事件通知核心(匹配与发送) | cmd/event-notification.go |
| 事件名枚举与复合展开 | internal/event/name.go |
| 通知配置(XML 规则模型) | internal/event/config.go |
| 前缀/后缀过滤器 | internal/event/rules.go |
| ARN 解析与校验 | internal/event/arn.go |
| 目标列表与并发发送 | internal/event/targetlist.go |
| S3 ListenNotification API | cmd/listen-notification-handlers.go |
按以上步骤,你可以为 MinIO 部署完整配置任意组合的桶事件通知:从 mc admin config set 注册端点,到 mc event add 绑定 ARN 与过滤器,再到用各目标的原生客户端(kafkacat、nsq_tail、redis-cli、curl、psql、mysql、Go/Python 订阅程序)完成端到端验证,并借助持久事件存储保证目标端离线期间事件不丢失。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0623
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00