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

基于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)的事件监听机制,无法实时感知呼叫结束(通道空闲)事件
  • 缺乏循环补呼逻辑,呼叫发起后不会自动补充新呼叫填满空闲通道
解决方案

修改后的代码实现以下核心功能:

  1. 监听AMI的Hangup事件,实时跟踪呼叫结束状态
  2. 维护活跃呼叫计数,定期检查可用通道数
  3. 循环运行,当活跃通道数小于并发限制时自动发起新呼叫,直到所有联系人处理完毕

修改后的代码:

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 09:42:03