背景 当社のビジネスは成長期にありましたが、サーバー側の性能はデモレベルに低い状態でした。そのため、サーバーのリファクタリングを開始しました。 テストの分野は多くのインターネット企業で失敗に終わりがちですが、幸運なことに、当社のCTOはテストを非常に重視しています。彼の口癖は「テスト不可能なプログラムは信頼できない」というもので、そのため当社のすべてのプログラムには対応するテストツールが用意されています。この記事で紹介するツールは、サーバー専用のテスト用に作成されたクライアントTSPです。
ビジネス要件 サーバーのリファクタリングにより、100万クライアントレベルの接続を実現します
実装原理
サーバーは3つのポートをリッスンし、これらのポートはすべてデバイスから開始されます tspはクライアントをシミュレートしてサーバーに接続し、デバイスとユーザーのオンラインを実現します プロセスプールにスレッドプールをネストする方式を採用し、高同時実行を実現します コード内でパケット送信時間の平均値、パケット受信時間の平均値、パケット受信バイトサイズの平均値、パケット受信エラーレートを統計します サーバーはハードウェア構成、CPU負荷、メモリ負荷、ディスクI/Oなどを収集します
具体的な実装
- メインフローの構築 - 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}')
- プロセスプールとスレッドプールの実装
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()})の実行完了')
- 統計パラメータの計算と出力
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のグローバルロックの制限によりスレッドにはいくつかの欠点がありますが、同時実行性の面ではプロセスよりも優れています。この方法で高密度のビジネスロジック操作を行うことで、サーバーのボトルネックを簡単に見つけることができます。