You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过Python的cx_Oracle CQN将消息发布到Amazon SNS

在Python中通过cx_Oracle CQN向Amazon SNS发送变更通知

1. 依赖与环境准备

  • 安装必要的Python包:
    pip install cx_Oracle boto3
    
  • 配置AWS凭证:可通过环境变量、~/.aws/credentials文件或IAM角色(运行在AWS服务上时)配置,确保拥有SNS主题的发布权限。
  • 确保Oracle数据库用户拥有CHANGE NOTIFICATION权限,执行授权语句:
    GRANT CHANGE NOTIFICATION TO your_db_user;
    

2. 完整实现代码

以下代码实现了CQN监听Oracle表的插入/更新操作,并将通知发送到指定SNS主题:

import cx_Oracle
import boto3
import os

# 初始化SNS客户端与主题ARN
sns_client = boto3.client('sns', region_name='your-aws-region')
SNS_TOPIC_ARN = 'arn:aws:sns:your-aws-region:your-account-id:your-topic-name'

# Oracle数据库连接配置(建议通过环境变量传递敏感信息)
DB_USER = os.getenv('ORACLE_USER')
DB_PASSWORD = os.getenv('ORACLE_PASSWORD')
DB_DSN = os.getenv('ORACLE_DSN')  # 格式示例:'db-host:1521/orcl'

def cqn_change_callback(message):
    """处理CQN通知并发送到SNS"""
    # 解析通知核心信息
    notification_data = {
        "event_type": message.type,
        "affected_tables": [f"{table.schema}.{table.name}" for table in message.tables],
        "notification_time": message.timestamp.strftime("%Y-%m-%d %H:%M:%S"),
        "query_id": message.queryId
    }

    # 若开启ROWID跟踪,附加变更行信息
    if message.rows:
        notification_data["changed_row_ids"] = [row.rowid for row in message.rows]

    # 发送到SNS主题
    try:
        sns_response = sns_client.publish(
            TopicArn=SNS_TOPIC_ARN,
            Message=str(notification_data),
            Subject="Oracle Table Change Notification"
        )
        print(f"通知已发送到SNS,消息ID: {sns_response['MessageId']}")
    except Exception as e:
        print(f"SNS发送失败: {str(e)}")

def start_cqn_listener():
    """启动CQN监听服务"""
    # 建立Oracle连接
    with cx_Oracle.connect(DB_USER, DB_PASSWORD, DB_DSN) as connection:
        # 注册CQN监听:指定目标表、监听事件类型、回调函数
        registration = connection.register(
            [('YOUR_SCHEMA', 'YOUR_TABLE')],  # 替换为需要监听的表(模式.表名)
            cx_Oracle.EVENT_DERIVED | cx_Oracle.EVENT_NOTIFY_ROWIDS,  # 监听衍生事件+返回ROWID
            callback=cqn_change_callback
        )

        print("CQN监听已启动,等待表变更...")
        try:
            # 持续监听,每5秒轮询一次
            while True:
                connection.waitfor(5)
        except KeyboardInterrupt:
            print("\n正在停止CQN监听...")
            registration.unregister()

if __name__ == "__main__":
    start_cqn_listener()

3. 关键配置说明

  • 监听事件调整:若需要监听DELETE操作,无需修改注册参数,cx_Oracle.EVENT_DERIVED已包含所有DML变更;若仅需监听特定操作,可替换为cx_Oracle.EVENT_INSERT/cx_Oracle.EVENT_UPDATE等。
  • ROWID跟踪:如果不需要获取具体变更行的ROWID,可移除cx_Oracle.EVENT_NOTIFY_ROWIDS参数,减少数据库资源消耗。
  • 多表监听:在register方法的第一个参数中添加多个表元组即可,例如[('SCHEMA1', 'TABLE1'), ('SCHEMA2', 'TABLE2')]。

内容的提问来源于stack exchange,提问作者Shivani

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.20 05:35:29