首页
/ MinIO 桶通知实战:把对象事件发布到 Kafka、RabbitMQ、Redis 等十类目标端

MinIO 桶通知实战:把对象事件发布到 Kafka、RabbitMQ、Redis 等十类目标端

2026-09-05 22:39:58作者:伍霜盼Ellen

本文以 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.gointernal/event/name.gointernal/event/arn.gocmd/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 枚举(如 ObjectCreatedPutObjectRemovedDeleteObjectReplicationFailedObjectRestorePost 等),且源码注释明确指出 s3:Replication:OperationCompletedReplication 属于 MinIO 扩展事件(非 AWS S3 原生语义)。Name 类型还提供 Expand() 方法,把复合类型展开为具体事件集:ObjectCreatedAll 展开为 CompleteMultipartUpload/Copy/Post/Put/PutRetention/PutLegalHold/PutTagging/DeleteTaggingObjectRemovedAll 展开为 Delete/DeleteMarkerCreated/NoOP/DeleteAllVersions——这正是 mc event adds3:ObjectCreated:* 这类通配写法能生效的底层机制。事件类型还采用位掩码设计(objectSingleTypesEnd 上限 64 个,用 uint64 位表示),保证 Match 匹配时的性能。

二、支持的通知目标总览

桶事件可以发布到以下十类目标,文档按目标逐一给出完整的三步配置法:

目标 目标 目标
AMQP(RabbitMQ) Redis MySQL
MQTT NATS Apache Kafka
Elasticsearch PostgreSQL Webhooks
NSQ

MinIO 客户端 mcevent 子命令用于设置与监听事件通知;各语言的 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>(如 1myinstance)作为标识符。

配置修改流程统一为: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.goparseARN 函数可见:字符串必须以 arn:minio:sqs: 开头、按冒号切分后必须恰好 6 段且第 5、6 段非空,否则返回 ErrInvalidARN;其中 <id> 对应配置名(如 1myinstance),<type> 对应子系统类型(如 amqpkafka)。这也是为什么自定义 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.SendbucketRulesMap[bucket].Match(eventName, objectName) 按桶规则+前缀/后缀过滤器算出目标集合 targetIDSet,无匹配则直接返回;
  • 事件体由 ToEvent 构造:EventVersion: "2.0"EventSource: "minio:s3"Sequencer 取对象修改时间的纳秒十六进制值,响应头包含 x-amz-request-idx-minio-origin-endpointx-minio-deployment-id 等 MinIO 自定义字段;删除类事件(Delete/DeleteMarkerCreated/NoOP)不携带 ETag/Size/ContentType
  • 是否同步发送由 API 配置 MINIO_API_SYNC_EVENTS 控制(见 SendisSyncEventsEnabled())。

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-bucketminio-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

brokertopic 必填。环境变量前缀 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.keysizesequencer 等字段)。

六、Elasticsearch

安装 Elasticsearch 后接入。此目标支持 namespaceaccess 两种格式:

  • 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

urlindexformat 必填,环境变量前缀 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

addresskeyformat 必填,环境变量前缀 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

addresssubject 必填;环境变量前缀 MINIO_NOTIFY_NATS_*,含 MINIO_NOTIFY_NATS_STREAMINGMINIO_NOTIFY_NATS_STREAMING_CLUSTER_IDMINIO_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.gostan.Connect("test-cluster", "test-client", stan.NatsURL(...), stan.SetConnectionLostHandler(...)) 建立连接并订阅 bucketevents,断线时在 SetConnectionLostHandler 回调中重连。收到的事件 JSON 结构与上面一致,且包含 contentTypeversionIduserDefined 等更完整的对象字段。

九、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_addresstopic 必填;环境变量前缀 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"}}]}

十四、实践要点与源码索引

  1. 过滤粒度mc event add 支持 --event(指定事件,如 put)、--prefix/--suffix(前缀/后缀过滤)。过滤器校验逻辑在 internal/event/rules.go;每个桶的规则集以 RulesMap(规则 → 目标集合)形式缓存在 EventNotifier.bucketRulesMap 中,热路径 Send 只需一次加读锁匹配。
  2. 多端点:任何目标都可通过 notify_xxx:<name> 注册多实例,事件会被同时推送到 ARN 引用的所有端点;internal/event/targetid.go 定义了 TargetID 及其到 ARN 字符串的互转。
  3. 可靠性分层:队列类目标(AMQP/MQTT/NATS/Kafka/NSQ/Webhook)与存储类目标(Redis/ES/PostgreSQL/MySQL)都可用 queue_dir/queue_limit 启用磁盘暂存,目标端恢复后自动重放;queue_limit 默认 100000 条。
  4. 语义差异:消息队列类目标是"事件流"(append-only);而 Redis/ES/PostgreSQL/MySQL 的 namespace 格式是"状态镜像"(upsert/delete 与桶内对象生命周期一致),access 格式则退化为"事件日志"。选型时按消费端是否需要全量当前状态来区分。
  5. 可观测性:事件体中的 x-amz-request-idx-minio-deployment-idx-minio-origin-endpoint 可与访问日志、复制指标交叉关联,便于审计链路排查(字段构造见 cmd/event-notification.goToEvent)。
  6. 相关文档与源码索引
主题 路径
本文档原文 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 订阅程序)完成端到端验证,并借助持久事件存储保证目标端离线期间事件不丢失。

登录后查看全文
热门项目推荐
相关项目推荐