Pythonを使用したサーバー高同時実行性能テストの実装

背景 当社のビジネスは成長期にありましたが、サーバー側の性能はデモレベルに低い状態でした。そのため、サーバーのリファクタリングを開始しました。 テストの分野は多くのインターネット企業で失敗に終わりがちですが、幸運なことに、当社のCTOはテストを非常に重視しています。彼の口癖は「テスト不可能なプログラムは信頼できない」というもので、そのため当社のすべてのプログラムには対応するテストツールが用意されています。この記事で紹介するツールは、サーバー専用のテスト用に作成されたクライアントTSPです。

ビジネス要件 サーバーのリファクタリングにより、100万クライアントレベルの接続を実現します

実装原理

サーバーは3つのポートをリッスンし、これらのポートはすべてデバイスから開始されます tspはクライアントをシミュレートしてサーバーに接続し、デバイスとユーザーのオンラインを実現します プロセスプールにスレッドプールをネストする方式を採用し、高同時実行を実現します コード内でパケット送信時間の平均値、パケット受信時間の平均値、パケット受信バイトサイズの平均値、パケット受信エラーレートを統計します サーバーはハードウェア構成、CPU負荷、メモリ負荷、ディスクI/Oなどを収集します

具体的な実装

  1. メインフローの構築 - socketベースの要求送受信
#!/usr/bin/env python
# -*- coding:utf-8 -*-

import socket
import struct
import logging

class NetworkConnector:
    def __init__(self, identifier, connection_params):
        self.connection_params = connection_params
        self.identifier = identifier
        self.__establish_connection(connection_params)
        
    def __establish_connection(self, address):
        try:
            logging.debug(f'{self.identifier} 接続データ: {address}')
            self.tcp_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
            self.tcp_socket.connect(address)
            self.tcp_socket.setblocking(1)
            self.tcp_socket.settimeout(1)
            logging.debug(f'{self.identifier} ソケット接続: 成功')
        except Exception as msg:
            logging.error(f'{self.identifier} ソケット接続結果: {msg}')
            
    def send_data(self, data):
        try:
            try:
                result = self.tcp_socket.send(data)
            except Exception as msg:
                logging.error(msg)
                raise msg
            logging.debug(f'{self.identifier} 送信データ長: {result}')
            return result
        except Exception as msg:
            logging.error(f'{self.identifier} 送信エラー: {msg}')
            raise
            
    def receive_data(self, size, format_code):
        result = ''
        try:
            try:
                result = self.tcp_socket.recv(int(size))
                while len(result) < int(size):
                    result += self.tcp_socket.recv(int(size) - len(result))
            except:
                pass
                
            if not result:
                logging.error(f'{self.identifier} 受信結果エラー: 結果がnullです')
                
            logging.debug(f'{self.identifier} 受信結果長: {len(result)}')
            data_structure = struct.unpack(format_code, result)
            logging.debug(f'{self.identifier} 受信結果展開: {data_structure}')
            return data_structure
        except Exception as msg:
            logging.error(f'{self.identifier} 受信エラー: {msg}')
            raise
            
    def receive_only(self, size):
        try:
            result = self.tcp_socket.recv(int(size))
            while len(result) < int(size):
                result += self.tcp_socket.recv(int(size) - len(result))
            logging.debug(f'{self.identifier} 受信結果長: {len(result)}')
            logging.debug(f'受信データ: {result}')
            return result
        except Exception as msg:
            logging.error(f'{self.identifier} 受信エラー: {msg}')
            raise
            
    def close_connection(self):
        try:
            logging.debug(f'{self.identifier} ソケットが閉じられました')
            self.tcp_socket.close()
        except Exception as msg:
            logging.debug(f'{self.identifier} ソケットクローズエラー: {msg}')
  1. プロセスプールとスレッドプールの実装
import threading
import multiprocessing
import os

def business_logic():
    print('ビジネスロジック実行中')

class CustomThread(threading.Thread):
    def __init__(self, function, arguments, name=''):
        threading.Thread.__init__(self)
        self.name = name
        self.function = function
        self.arguments = arguments
        
    def run(self):
        self.function(*self.arguments)

def thread_manager():
    global sequence_number
    threads = []
    
    for i in range(thread_count):  # thread_countは同時実行スレッド数
        mac, mac_real, sequence_number = get_device_info()
        t = CustomThread(business_logic, (mac, mac_real, sequence_number))
        threads.append(t)
        
    for i in range(thread_count):
        threads[i].start()
        
    for i in range(thread_count):
        threads[i].join()

if __name__ == '__main__':
    result = ''
    process_pool = multiprocessing.Pool(processes=process_count)  # process_countはプロセスプールの数
    
    logging.info(f"メインプロセス({os.getpid()})を実行中...")
    
    for i in range(process_num):  # process_numは同時実行プロセス数
        result = process_pool.apply_async(thread_manager)
        
    process_pool.close()
    process_pool.join()
    
    if not result.successful():
        logging.error(f'メインプロセスの異常: {result.successful()}')
    else:
        logging.info(f'終了: メインプロセス({os.getpid()})の実行完了')
  1. 統計パラメータの計算と出力
import json

class GlobalVariableManager:
    # 辞書としてグローバル変数を構築 - mapを参考に変数の追加、削除、変更、検索を含む
    variable_map = {}
    
    def set_variable(self, key, value):
        if isinstance(value, dict):
            value = json.dumps(value)
        self.variable_map[key] = value
        logging.debug(f"{key}: {value}")
        
    def set_multiple(self, **kwargs):
        try:
            for key, value in kwargs.items():
                self.variable_map[key] = str(value)
                logging.debug(f"{key}: {value}")
        except Exception as msg:
            logging.error(msg)
            raise
            
    def delete_variable(self, key):
        try:
            del self.variable_map[key]
            return self.variable_map
        except KeyError:
            logging.error(f"キー '{key}' が存在しません")
            
    def get_variables(self, *args):
        try:
            result = {}
            for key in args:
                if len(args) == 1:
                    result = self.variable_map[key]
                    logging.debug(f"{key}: {result}")
                elif len(args) == 1 and args[0] == 'all':
                    result = self.variable_map
                else:
                    result[key] = self.variable_map[key]
            return result
        except KeyError:
            logging.warning(f"キー '{key}' が存在しません")
            return 'Null_'

まとめ マルチプロセス方式とマルチスレッド方式を比較すると、Pythonのグローバルロックの制限によりスレッドにはいくつかの欠点がありますが、同時実行性の面ではプロセスよりも優れています。この方法で高密度のビジネスロジック操作を行うことで、サーバーのボトルネックを簡単に見つけることができます。

タグ: Python サーバー性能テスト 高同時実行 マルチスレッド マルチプロセス

8月5日 19:51 投稿