如何通过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
相关产品推荐
相关产品推荐

