基于Python的Asterisk AMI批量并发外呼实现及代码优化需求
问题描述
我需要从数据库中同时给一组联系人(假设组内有1000个联系人)发起外呼。例如我有30条并发通道,要求代码同时向30个号码发起呼叫,当有通道空闲时,自动补呼对应数量的号码,持续此流程。
当前实现的代码存在问题:它会逐个发起呼叫,只有当一个呼叫被接听或挂断后,才会发起下一个呼叫,这不符合我的需求。
原代码:
import time import asterisk.manager import pymysql # Assuming you have pymysql library installed, you can install it with: pip install pymysql DB_HOST = '127.0.0.1' DB_USER = 'root' DB_PASSWORD = 'XXXX' DB_NAME = 'abcd' AMI_USERNAME = 'admin' AMI_PASSWORD = 'XYZ' AMI_HOST = '127.0.0.1' AMI_PORT = 5038 TOTAL_CHANNELS_LIMIT = 30 # Set your total channel limit here def connect_to_db(): connection = pymysql.connect(host=DB_HOST, user=DB_USER, password=DB_PASSWORD, db=DB_NAME, cursorclass=pymysql.cursors.DictCursor) return connection def connect_to_ami(): manager = asterisk.manager.Manager() manager.connect(AMI_HOST, AMI_PORT) manager.login(AMI_USERNAME, AMI_PASSWORD) return manager def check_active_channels(manager): response = manager.command('core show channels concise') active_channels = [line.split()[0] for line in response.data.split('\n') if line.strip()] return active_channels def get_active_campaigns(connection): with connection.cursor() as cursor: cursor.execute("SELECT name as campaign_name, call_group, fixed_channels FROM vb_schedule_play WHERE status = 'active'") result = cursor.fetchall() active_campaigns = [{'name': row['campaign_name'], 'fixed_channels': row['fixed_channels'], 'call_group': row['call_group']} for row in result] return active_campaigns def get_contacts_for_campaign(connection, campaign_name): with connection.cursor() as cursor: sql = f"SELECT phone as contact_number FROM contact_list WHERE group_name = 'sohub_group' AND (status IS NULL OR status = '') LIMIT 1" cursor.execute(sql) result = cursor.fetchone() if result: return result['contact_number'] else: return None def update_contact_status(connection, contact_number, campaign_name, status): with connection.cursor() as cursor: sql = f"UPDATE contact_list SET status = '{status}' WHERE phone = '{contact_number}'" cursor.execute(sql) connection.commit() def initiate_calls(manager, campaign_name, fixed_channels, connection, total_channels_limit=TOTAL_CHANNELS_LIMIT): if fixed_channels is None: fixed_channels = 0 fixed_channels = int(fixed_channels) if fixed_channels is not None else 0 if fixed_channels > 0: # Use fixed_channels for the specified campaign channels_to_use = min(fixed_channels, total_channels_limit) else: # Dynamic channel allocation for campaigns with fixed_channels equal to 0 available_channels = check_active_channels(manager) channels_to_use = min(len(available_channels), total_channels_limit) total_channels_used = 0 # Initialize the total channels used for _ in range(channels_to_use): if total_channels_used >= total_channels_limit: print(f"Reached the total channel limit ({total_channels_limit}). Stopping further calls.") break contact_number = get_contacts_for_campaign(connection, campaign_name) if contact_number: full_channel = f'PJSIP/{contact_number}' manager.originate( channel=full_channel, exten='s', context='ami-contact-center', priority=1, caller_id='0123456789', variables={ 'CallerID': '0123456789', 'CALLERID(all)': '0123456789', 'play_type': 'tts', 'play': 'hello this is a test call', }, timeout=30000 ) update_contact_status(connection, contact_number, campaign_name, 'calling') total_channels_used += 1 # Increment the total channels used time.sleep(1) # Add a delay to avoid overwhelming the system else: print(f"No more contacts available for campaign {campaign_name}") def main(): db_connection = connect_to_db() ami_manager = connect_to_ami() active_campaigns = get_active_campaigns(db_connection) for campaign in active_campaigns: campaign_name = campaign['name'] fixed_channels = campaign['fixed_channels'] initiate_calls(ami_manager, campaign_name, fixed_channels, db_connection, total_channels_limit=TOTAL_CHANNELS_LIMIT) ami_manager.logoff() db_connection.close() if __name__ == '__main__': main()
问题分析
原代码的核心缺陷:
- 仅在启动时一次性发起固定数量的呼叫,没有持续监控通道状态的逻辑
- 未利用Asterisk Manager Interface(AMI)的事件监听机制,无法实时感知呼叫结束(通道空闲)事件
- 缺乏循环补呼逻辑,呼叫发起后不会自动补充新呼叫填满空闲通道
解决方案
修改后的代码实现以下核心功能:
- 监听AMI的
Hangup事件,实时跟踪呼叫结束状态 - 维护活跃呼叫计数,定期检查可用通道数
- 循环运行,当活跃通道数小于并发限制时自动发起新呼叫,直到所有联系人处理完毕
修改后的代码:
import time import threading import asterisk.manager import pymysql DB_HOST = '127.0.0.1' DB_USER = 'root' DB_PASSWORD = 'XXXX' DB_NAME = 'abcd' AMI_USERNAME = 'admin' AMI_PASSWORD = 'XYZ' AMI_HOST = '127.0.0.1' AMI_PORT = 5038 TOTAL_CHANNELS_LIMIT = 30 active_calls = 0 lock = threading.Lock() stop_flag = False def connect_to_db(): return pymysql.connect( host=DB_HOST, user=DB_USER, password=DB_PASSWORD, db=DB_NAME, cursorclass=pymysql.cursors.DictCursor ) def connect_to_ami(): manager = asterisk.manager.Manager() manager.connect(AMI_HOST, AMI_PORT) manager.login(AMI_USERNAME, AMI_PASSWORD) # 注册Hangup事件监听,实时更新活跃呼叫数 manager.register_event('Hangup', handle_hangup) return manager def check_active_channels(manager): try: response = manager.command('core show channels concise') active_channels = [line.split()[0] for line in response.data.split('\n') if line.strip()] return len(active_channels) except Exception as e: print(f"检查活跃通道出错: {e}") return 0 def get_active_campaigns(connection): with connection.cursor() as cursor: cursor.execute("SELECT name as campaign_name, call_group, fixed_channels FROM vb_schedule_play WHERE status = 'active'") result = cursor.fetchall() return [ { 'name': row['campaign_name'], 'fixed_channels': int(row['fixed_channels']) if row['fixed_channels'] else 0, 'call_group': row['call_group'] } for row in result ] def get_next_contact(connection, campaign_name, call_group): with connection.cursor() as cursor: # 使用FOR UPDATE SKIP LOCKED避免多线程重复获取同一联系人 sql = """ SELECT phone as contact_number FROM contact_list WHERE group_name = %s AND (status IS NULL OR status = '') LIMIT 1 FOR UPDATE SKIP LOCKED """ cursor.execute(sql, (call_group,)) result = cursor.fetchone() if result: contact = result['contact_number'] # 先标记为呼叫中,防止重复获取 update_contact_status(connection, contact, campaign_name, 'calling') return contact return None def update_contact_status(connection, contact_number, campaign_name, status): try: with connection.cursor() as cursor: sql = "UPDATE contact_list SET status = %s WHERE phone = %s" cursor.execute(sql, (status, contact_number)) connection.commit() except Exception as e: print(f"更新联系人状态出错: {e}") def handle_hangup(event, manager): global active_calls with lock: active_calls -= 1 print(f"呼叫结束,当前活跃呼叫数: {active_calls}") def initiate_call(manager, contact_number, campaign_name): global active_calls try: full_channel = f'PJSIP/{contact_number}' # 使用异步发起呼叫,避免阻塞线程 manager.originate( channel=full_channel, exten='s', context='ami-contact-center', priority=1, caller_id='0123456789', variables={ 'CallerID': '0123456789', 'CALLERID(all)': '0123456789', 'play_type': 'tts', 'play': 'hello this is a test call', }, timeout=30000, async_=True ) with lock: active_calls += 1 print(f"已发起呼叫至 {contact_number},当前活跃呼叫数: {active_calls}") return True except Exception as e: print(f"发起呼叫至 {contact_number} 失败: {e}") # 发起失败,重置联系人状态 db_conn = connect_to_db() update_contact_status(db_conn, contact_number, campaign_name, '') db_conn.close() return False def run_campaign(manager, campaign): db_conn = connect_to_db() campaign_name = campaign['name'] call_group = campaign['call_group'] # 计算该campaign可用的最大通道数 campaign_channel_limit = min(campaign['fixed_channels'], TOTAL_CHANNELS_LIMIT) if campaign['fixed_channels'] > 0 else TOTAL_CHANNELS_LIMIT while not stop_flag: with lock: current_active = active_calls # 计算可发起的新呼叫数量 available_slots = campaign_channel_limit - current_active if available_slots <= 0: time.sleep(1) continue for _ in range(available_slots): contact = get_next_contact(db_conn, campaign_name, call_group) if not contact: print(f"campaign {campaign_name} 已无可用联系人") global stop_flag stop_flag = True break initiate_call(manager, contact, campaign_name) # 避免短时间内发起过多呼叫 time.sleep(0.5) time.sleep(1) db_conn.close() def main(): global stop_flag try: db_connection = connect_to_db() ami_manager = connect_to_ami() active_campaigns = get_active_campaigns(db_connection) if not active_campaigns: print("未找到活跃的外呼任务") return # 为每个campaign启动独立线程 threads = [] for campaign in active_campaigns: thread = threading.Thread(target=run_campaign, args=(ami_manager, campaign)) threads.append(thread) thread.start() # 等待所有线程执行完毕 for thread in threads: thread.join() except Exception as e: print(f"主程序出错: {e}") stop_flag = True finally: ami_manager.logoff() db_connection.close() if __name__ == '__main__': main()
关键改进点
- 事件驱动监控:注册
Hangup事件处理器,实时更新活跃呼叫计数 - 异步呼叫发起:使用
async_=True参数避免阻塞主线程 - 线程化任务:每个外呼任务独立运行线程,支持多任务并发
- 数据库锁机制:
SELECT ... FOR UPDATE SKIP LOCKED防止多线程重复获取同一联系人 - 循环补呼逻辑:持续检查可用通道数,自动补充新呼叫直到所有联系人处理完成
内容的提问来源于stack exchange,提问作者Tanvir Islam
相关产品推荐
相关产品推荐

