使用只读路由和可用性组

Always On可用性组为SQL Server数据库提供高可用性。 在连接之前,数据库管理员必须 创建一个可用性组监听器 ,并在服务器 上配置只读路由 。 mssql-python 驱动支持通过以下连接字符串关键字进行可用性组连接:

  • ApplicationIntent - 将连接路由到只读的二级副本
  • MultiSubnetFailover - 启用跨子网的并行连接尝试以加快故障切换

连接可用性组监听器

基本侦听器连接

通过监听器DNS名称连接到可用性组,而不是特定服务器:

import mssql_python

# Connect via availability group listener
conn = mssql_python.connect(
    "Server=ag-listener.contoso.com;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryDefault;"
    "Encrypt=yes;"
)

cursor = conn.cursor()
cursor.execute("SELECT @@SERVERNAME AS ServerName")
print(f"Connected to: {cursor.fetchval()}")

带有端口规格

当监听者使用非默认端口时,指定端口号:

conn = mssql_python.connect(
    "Server=ag-listener.contoso.com,1433;"  # Listener with port
    "Database=<database>;"
    "Authentication=ActiveDirectoryDefault;"
    "Encrypt=yes;"
)

只读路由

启用只读意向

使用 ApplicationIntent=ReadOnly 路由到辅助副本:

# Connect for read operations - routes to secondary
read_conn = mssql_python.connect(
    "Server=ag-listener.contoso.com;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryDefault;"
    "ApplicationIntent=ReadOnly;"
    "Encrypt=yes;"
)

# Connect for read-write operations - routes to primary
write_conn = mssql_python.connect(
    "Server=ag-listener.contoso.com;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryDefault;"
    "ApplicationIntent=ReadWrite;"  # Default
    "Encrypt=yes;"
)

验证路由

确认是哪个副本处理了连接,以及它是主副本还是备份副本:

def check_replica_role(conn) -> str:
    """Check if connected to primary or secondary."""
    cursor = conn.cursor()
    cursor.execute("""
        SELECT 
            @@SERVERNAME AS ServerName,
            CASE 
                WHEN DATABASEPROPERTYEX(DB_NAME(), 'Updateability') = 'READ_WRITE' 
                THEN 'Primary'
                ELSE 'Secondary'
            END AS Role
    """)
    row = cursor.fetchone()
    return f"{row.ServerName} ({row.Role})"

read_conn = mssql_python.connect(connection_string + "ApplicationIntent=ReadOnly;")
print(f"Read connection: {check_replica_role(read_conn)}")

write_conn = mssql_python.connect(connection_string + "ApplicationIntent=ReadWrite;")
print(f"Write connection: {check_replica_role(write_conn)}")

多子网故障转移

启用多子网故障切换

对于跨多个子网的可用性组:

conn = mssql_python.connect(
    "Server=ag-listener.contoso.com;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryDefault;"
    "MultiSubnetFailover=yes;"
    "Encrypt=yes;"
)

此设置:

  • 尝试并行连接所有IP地址。
  • 在多子网配置中减少故障切换时间。
  • 最适用于所有可用性组连接。

连接架构模式

独立的读写连接

使用单独的连接,将读流量发送到辅助节点,并将写入流量发送到主节点:

class DatabaseConnections:
    """Manage separate connections for read and write operations."""
    
    def __init__(self, listener: str, database: str):
        self.base_conn_str = f"Server={listener};Database={database};Authentication=ActiveDirectoryDefault;Encrypt=yes;"
        self._read_conn = None
        self._write_conn = None
    
    @property
    def read_connection(self):
        """Get or create read-only connection (secondary replica)."""
        if self._read_conn is None:
            self._read_conn = mssql_python.connect(
                self.base_conn_str + "ApplicationIntent=ReadOnly;MultiSubnetFailover=yes;"
            )
        return self._read_conn
    
    @property
    def write_connection(self):
        """Get or create read-write connection (primary replica)."""
        if self._write_conn is None:
            self._write_conn = mssql_python.connect(
                self.base_conn_str + "ApplicationIntent=ReadWrite;MultiSubnetFailover=yes;"
            )
        return self._write_conn
    
    def close(self):
        if self._read_conn:
            self._read_conn.close()
        if self._write_conn:
            self._write_conn.close()

# Usage
db = DatabaseConnections("ag-listener.contoso.com", "<database>")

# Queries go to secondary
cursor = db.read_connection.cursor()
cursor.execute("SELECT ProductID, Name, ListPrice FROM Production.Product")
products = cursor.fetchall()

# Writes go to primary
cursor = db.write_connection.cursor()
cursor.execute("INSERT INTO #Products (Name) VALUES (%(name)s)", {"name": "New Product"})
db.write_connection.commit()

db.close()

读后写模式

先写入主节点,然后在次节点上验证该数据,并将复制延迟考虑在内:

def create_order_and_verify(conn_manager, order_data: dict):
    """Create order on primary, verify on secondary with eventual consistency."""
    
    # Write to primary
    write_cursor = conn_manager.write_connection.cursor()
    write_cursor.execute("""
        INSERT INTO #Orders (CustomerID, Total) VALUES (%(cust)s, %(total)s);
        SELECT SCOPE_IDENTITY();
    """, order_data)
    order_id = write_cursor.fetchval()
    conn_manager.write_connection.commit()
    
    # Wait for replication (in production, use more sophisticated approach)
    import time
    time.sleep(1)
    
    # Verify on secondary
    read_cursor = conn_manager.read_connection.cursor()
    read_cursor.execute("SELECT * FROM #Orders WHERE OrderID = %(id)s", {"id": order_id})
    
    if read_cursor.fetchone():
        print(f"Order {order_id} replicated to secondary")
    else:
        print(f"Order {order_id} not yet replicated")
    
    return order_id

处理故障转移

连接弹性

当故障切换中断连接时,自动重试查询:

import time

def execute_with_failover_retry(conn_str: str, query: str, params: dict,
                                max_retries: int = 3) -> list:
    """Execute query with automatic reconnection on failover."""
    
    for attempt in range(max_retries + 1):
        try:
            conn = mssql_python.connect(conn_str)
            cursor = conn.cursor()
            cursor.execute(query, params)
            results = cursor.fetchall()
            conn.close()
            return results
        except mssql_python.OperationalError as e:
            error_str = str(e)
            # Check for failover-related errors
            # Azure SQL transient error codes
            # See: /azure/azure-sql/database/troubleshoot-common-errors-issues
            if any(code in error_str for code in ["40613", "40197", "40501"]):
                if attempt < max_retries:
                    print(f"Failover detected, retrying ({attempt + 1}/{max_retries})...")
                    time.sleep(5 * (attempt + 1))  # Exponential backoff
                    continue
            raise

# Usage
results = execute_with_failover_retry(
    "Server=<ag-listener>.contoso.com;Database=<database>;Authentication=ActiveDirectoryDefault;MultiSubnetFailover=yes;Encrypt=yes;",
    "SELECT ProductID, Name, ListPrice FROM Production.Product WHERE ProductSubcategoryID = %(cat)s",
    {"cat": 5}
)

重试循环仅对瞬态错误码做出反应。 非暂时性失败,如认证错误、权限错误或查询语法错误,会立即被触发,因为重试无法解决。

检测主要角色变化

检查当前的副本角色,如果主节点移动了,请重新连接:

def is_primary(conn) -> bool:
    """Check if current connection is to primary replica."""
    cursor = conn.cursor()
    cursor.execute("""
        SELECT DATABASEPROPERTYEX(DB_NAME(), 'Updateability') AS Updateability
    """)
    return cursor.fetchval() == 'READ_WRITE'

def ensure_primary(conn_str: str) -> mssql_python.Connection:
    """Ensure connection is to primary, reconnect if needed."""
    conn = mssql_python.connect(conn_str + "ApplicationIntent=ReadWrite;")
    
    if not is_primary(conn):
        # Might happen during failover
        conn.close()
        time.sleep(2)
        conn = mssql_python.connect(conn_str + "ApplicationIntent=ReadWrite;")
    
    return conn

报告工作负载

将报告转移到辅助端

对次要副本运行报告查询以减少主副本的负载:

class ReportingService:
    """Service that runs reports on secondary replicas."""
    
    def __init__(self, listener: str, database: str):
        self.conn_str = (
            f"Server={listener};"
            f"Database={database};"
            "Authentication=ActiveDirectoryDefault;"
            "ApplicationIntent=ReadOnly;"
            "MultiSubnetFailover=yes;"
            "Encrypt=yes;"
        )
    
    def run_report(self, report_query: str, params: dict = None) -> list:
        """Execute report query on secondary replica."""
        conn = mssql_python.connect(self.conn_str)
        cursor = conn.cursor()
        
        try:
            cursor.execute(report_query, params or {})
            return cursor.fetchall()
        finally:
            cursor.close()
            conn.close()
    
    def get_sales_summary(self, start_date, end_date) -> dict:
        """Run sales summary report."""
        results = self.run_report("""
            SELECT 
                YEAR(OrderDate) AS Year,
                MONTH(OrderDate) AS Month,
                COUNT(*) AS OrderCount,
                SUM(TotalDue) AS TotalSales
            FROM Sales.SalesOrderHeader
            WHERE OrderDate BETWEEN %(start)s AND %(end)s
            GROUP BY YEAR(OrderDate), MONTH(OrderDate)
            ORDER BY Year, Month
        """, {"start": start_date, "end": end_date})
        
        return [
            {"year": r.Year, "month": r.Month, 
             "orders": r.OrderCount, "sales": r.TotalSales}
            for r in results
        ]

# Usage
reports = ReportingService("ag-listener.contoso.com", "sales_db")
summary = reports.get_sales_summary(date(2024, 1, 1), date(2024, 12, 31))

最佳做法

始终使用 MultiSubnetFailover

在每个可用性组连接字符串中都包含 MultiSubnetFailover=yes,以实现更快的故障转移:

# Recommended for all AG connections
conn = mssql_python.connect(
    "Server=ag-listener.contoso.com;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryDefault;"
    "MultiSubnetFailover=yes;"  # Always include this
    "Encrypt=yes;"
)

将意图与操作匹配

根据每项操作是读取还是写入数据,将其路由到正确的副本:

def get_appropriate_connection(operation: str, connections: DatabaseConnections):
    """Get connection appropriate for the operation type."""
    read_operations = {"SELECT", "REPORT", "EXPORT", "ANALYTICS"}
    
    if operation.upper() in read_operations:
        return connections.read_connection
    else:
        return connections.write_connection

处理暂时不可用的辅助项

遇到瞬态错误时重试次要副本,如果次要副本无法访问,再退回主副本。

import time

def query_with_fallback(conn_manager, query: str, params: dict, max_retries: int = 2):
    """Query the secondary, retrying transient errors before falling back to the primary."""
    transient_codes = ["40613", "40197", "40501"]
    for attempt in range(max_retries + 1):
        try:
            cursor = conn_manager.read_connection.cursor()
            cursor.execute(query, params)
            return cursor.fetchall()
        except mssql_python.OperationalError as e:
            # Re-raise failures that retrying can't fix, such as authentication,
            # permission, or query syntax errors.
            if not any(code in str(e) for code in transient_codes):
                raise
            # Retry the secondary for transient errors before giving up.
            if attempt < max_retries:
                time.sleep(2 * (attempt + 1))
                continue
            # Secondary still unavailable after retries: fall back to the primary.
            print("Secondary unavailable, using primary")
            cursor = conn_manager.write_connection.cursor()
            cursor.execute(query, params)
            return cursor.fetchall()