前言经过之前几周的学习和实践项目的基础通讯骨架已经具有雏形——一个跨三端的具有各类基本组件拥有简单业务场景的网络模拟架构已经形成——但距离项目预设中的高拟真目标依旧很远Linux-serverAndroid-client和Windows-gateway并没有让之前的动态配置策略在多实例条件下被有效运用MySQL作为database没有经过合理的优化业务场景中的各类message在存储和传输中没有数据安全措施......总得来说任务很多。本篇会依据部分问题通过MySQL优化和各类组件的改造解决它们并给出完整的测试方案让项目的拟真水准更进一步。TIPS实践中出现的一些问题由于项目由我个人完成而我并非专业的测试人员很多时候项目的各类组件会因为修改/报错产生异步关闭Windows-gateway上的flask网关可能会因为修改而暂停但Android-client上的OKHTTP组件依然在尝试寻找那个已经关闭的端口这种操作频繁如果频繁的出现就会带来问题当前的Android-client和Windows-gateway通讯依赖于端口之间TCP的反代reverse这种链接会因为Windows系统中各种应用同时大量运行导致端口特别是进程端口的不固定出现坏死旧进程可能会占据5000端口新进程启动后没有接管portclient请求对象的错位进程卡死在RPC和lock中......总之比较头疼在之前的文章中我习惯把它叫做僵尸往往导致代码正常但是流程在某些地方卡死。下面给出相应的排查方案Get-NetTCPConnection -State Listen -LocalPort 5000 ##查看监听5000端口的port Get-Process -Id xxx ##根据上方指令的port获取相关进程 ##如果发现多个python/flask进程基本可以确定是旧进程问题如果一切正常但仍旧无法让流程跑通可以试试下一步防止是ADB reverse残留导致adb kill-server ##终止当前反代服务 adb start-server ##重新开始 adb connect 127.0.0.1:5555 adb -s 127.0.0.1:5555 reverse tcp:5000 tcp:5000 adb reverse --list或者是flask调用旧的python子进程use_reloaderFalse ##让状态重置然后重启服务通过上方排查大概率状态将会重置而未完成状态的进程也会被释放最极端情况可以按下方顺序关闭各个组件后重启Windows系统彻底重置虽然当前项目未涉及但是注册表相关报错上述措施都是无效的不要套用导致更大的链式错误类似的情况也可能发生在Linux-server和RabbitMQflask和RabbitMQ上只不过由于RabbitMQ作为中间件的自启动模式以及它相对固定的localIP和Port我目前并没有遇到严重的报错但我仍然建议你在测试时不要频繁单独关闭某一个组件按client-gateway-server顺序关闭时个比较好的选项如果遇到类似情况可以试试上文的排查方案。一.围绕MySQL的优化尝试当前的众多应用中数据库一般来说是需要重点关注的对象大多数情况下在一个完整的业务流程总是绕不开各种数据的存取其中还有相当多的一部分需要同步进行数据的二次加工。而database中的data——可过期的如LFU,LRU逻辑临时文档关系到核心业务的key/confirm_code用户的password乃至现实中的各类隐私——更需要妥善保存。同时当应用的data规模到了一定量级如何在保证安全的同时高效查找也是个问题很明显在之前的基础骨架中这些东西并没有得到妥善处理因此作为上述的类似问题的中心一个围绕MySQL的优化是必要的也是项目后续更多功能拓展的优化基础。1.更贴合现实的token功能扩展在之前的验证码权限升级场景中我们引用了token机制作为权限upgrade以及role切换的依据但是token的协议仅仅是为了快速验证功能——旧协议当中token从本质上来说和其它的message并没有本质的区别它们都在同一个packet中通过明文传输和存取Windows-gateway的检测机制也停留在对token的首位检验MySQL也并没有让token在我们预设的场景中token是和当前常见的API-key类似唯一有效身份凭证得到有效的隔离和索引——简单来说token是”裸奔“的在现实中这是不可接受的。下面我会结合相关的代码给出相应的解决方案。CREATE TABLE IF NOT EXISTS access_grant ( token_hash BINARY(32) NOT NULL, ##token的hash化储存 producer_id VARCHAR(64) NOT NULL, ##区分多实例的简单标识 session_id VARCHAR(64) NOT NULL, ##会话ID用于同一user/producer的不同role/level场景 role VARCHAR(32) NOT NULL, issued_at DATETIME(6) NOT NULL, expires_at DATETIME(6) NOT NULL, revoked_at DATETIME(6) NULL, last_used_at DATETIME(6) NULL, PRIMARY KEY (token_hash), ##hash化的token作为內键从而作为权威索引 INDEX idx_producer_expires ( ##过期/producer联合索引 producer_id, expires_at ), INDEX idx_expires_at (expires_at), ##单过期索引 INDEX idx_role_expires ( ##role/过期联合索引 role, expires_at ) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COLLATEutf8mb4_unicode_ci;在这个新建表逻辑中关键修改在于token不会用明文出现而是通过hash转化为更安全的形式通过token_hash和role的关联role和token_hash实现了有效的索引关联让关键的upgrade参数可以被Windows-gateway快速检测而不是依靠旧协议让gateway承担依照明文token算消去”裸奔“这种不优雅的做法。下面我会结合Windows-gateway代码修改让逻辑更加清晰def hash_access_token(token: str) - bytes: return hmac.new( TOKEN_HASH_PEPPER.encode(utf-8), token.encode(utf-8), hashlib.sha256, ).digest() ##通过使用SHA256的hash模式让原先string类型的token转成hash格式 ##SHA256也方便我们后续进行可能的微调比如调节盐度等但那是信息安全专项优化本篇不涉及token_hash之后会通过RabbitMQ发送到consumer中——注意当前的token_hash其实是base64模式后续的consumer需要解码最后才是通过consumer存储到MySQL中作为role/level的升级依据——这样就能够让旧协议的明文查找首位模式得到一次升级。def rabbit_rpc(message_type: str, payload: dict): ##RPC类型联合RabbitMQ做token_hash的收发 connection None try: ##基本的channel和connection设置联合pika实现 ...... reply channel.queue_declare( ##声名临时reply_queue方便数据收发 queue, ##无需频繁改动配置 exclusiveTrue, ) reply_queue reply.method.queue request_id str(uuid.uuid4()) ##随机ID生成 envelope { ##定义信封格式 message_id: request_id, source: flask_gateway, type: message_type, created_at: datetime.now(timezone.utc).isoformat(), payload: payload, } body json.dumps( ##信封读取 envelope, ensure_asciiFalse, ).encode(utf-8) properties pika.BasicProperties( ##权限设置 content_typeapplication/json, content_encodingutf-8, delivery_mode2, correlation_idrequest_id, reply_toreply_queue, ) channel.confirm_delivery() channel.basic_publish( ##发布 exchange, routing_keyQUEUE_NAME, bodybody, propertiesproperties, mandatoryTrue, ) deadline time.monotonic() RPC_TIMEOUT_SECONDS ##过期时效定义 while time.monotonic() deadline: method, response_properties, response_body ( ##在过期前发布 channel.basic_get( queuereply_queue, auto_ackTrue, ) ) if method is None: ##短暂休眠防止运行时的无效尝试 time.sleep(0.05) continue if response_properties.correlation_id ! request_id: continue return json.loads( response_body.decode(utf-8) ) return None finally: if connection is not None and connection.is_open: connection.close() def store_access_grant_rpc( ##通过RabbitMQ向consumer发布以token为中心的各种访问参数 access_token: str, producer_id: str, session_id: str, role: str, expires_at: float, ): token_hash hash_access_token(access_token) ##获取hash化的token issued_at datetime.now( ##生成日志 timezone.utc ).isoformat() expires_at_text datetime.fromtimestamp( ##过期日志 expires_at, timezone.utc, ).isoformat() return rabbit_rpc( TYPE_ACCESS_GRANT_STORE, { producer_id: producer_id, session_id: session_id, token_hash_b64: base64.b64encode( ##解码hashbase64格式 token_hash ).decode(ascii), role: role, issued_at: issued_at, expires_at: expires_at_text, }, ) def resolve_access_grant_rpc( ##token代表的各类参数 producer_id: str, access_token: str, ): token_hash hash_access_token(access_token) return rabbit_rpc( ##向rabbitmq发布的grant构造 TYPE_ACCESS_GRANT_RESOLVE, { producer_id: producer_id, token_hash_b64: base64.b64encode( token_hash ).decode(ascii), }, )通过这一部分代码应该能够让你了解token_hash的具体构造以及各类关键访问参数和token_hash的关系但需要注意的是在具体实现中不要把Windows-gateway和Linux-server的界限混淆比如尝试在gateway中通过mysql-connector直接操作MySQL这种越界行为是不可行的当前的gateway无法和Linux虚拟机上的MySQL直接交互同时也会违反gateway作为中间件的权限至少在当前Linux-server对应的consumer是唯一能够操作MySQL的组件。同时这里也涉及到了RFC类型SHA_hash等额外内容在后续的文章中应该会有一篇专门讲解这些内容顺带填上之前文章中类似的坑。consumer相关的代码我这里不会单独给出Linux-server在下文会做一些别的改造但根据上文的逻辑consumer应该可以参考往期文章做快速实现。2.MySQL的自更新在之前的文章中动态配置局限在role/level这种参数的快速改变producer_id也是修改一个JSON变量就能解决的但现实中的动态往往不仅是参数的改造更多的时候组件本身也会变动以当前的MySQL来举例就算我们将查询做到比较优秀的程度表的单一仍会导致索引结构的失效——一个表里面如果不断累积大量数据那么索引就会被数量拖垮——所以我选择引入一个storage_manager组件该组件将会根据系统时间来自动的创建新的表结构从而让MySQL的索引结构能够长久有效前提是表结构相对确定不大改。下面结合部分代码我会通过其中注释讲解具体思路DB_CONFIG { ##MySQL的具体配置这类配置一般不会轻易改动省略} PRE_CREATE_MONTHS int( ##按月份来作为自更新依据 os.environ.get( STORAGE_PRE_CREATE_MONTHS, 2, ) ) RETENTION_MONTHS int( ##12个月份为一个大周期 os.environ.get( STORAGE_RETENTION_MONTHS, 12, ) ) MANAGER_INTERVAL_SECONDS int( ##更新缓冲时间 os.environ.get( STORAGE_MANAGER_INTERVAL_SECONDS, 3600, ) ) SCHEMA_VERSION 1 ##版本可更新 CREATE_REGISTRY_SQL CREATE TABLE IF NOT EXISTS storage_registry ( ##建月份表逻辑作为月份总表 bucket_key CHAR(6) NOT NULL, table_name VARCHAR(64) NOT NULL, start_at DATETIME(6) NOT NULL, end_at DATETIME(6) NOT NULL, status VARCHAR(16) NOT NULL, schema_version INT NOT NULL, created_at DATETIME(6) NOT NULL, updated_at DATETIME(6) NOT NULL, PRIMARY KEY (bucket_key), UNIQUE KEY uk_storage_table_name (table_name), INDEX idx_status_start_at ( status, start_at ) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COLLATEutf8mb4_unicode_ci def utc_now(): ##当前时间 return datetime.now( timezone.utc ).replace(tzinfoNone) def month_start(value: datetime) - datetime: ##月份初始化 return value.replace( day1, hour0, minute0, second0, microsecond0, ) def add_months( ##新增月份 value: datetime, months: int, ) - datetime: month_index ( value.year * 12 value.month - 1 months ) year month_index // 12 ##年份计算 month month_index % 12 1 ##月份计算 return value.replace( ##年月值更新 yearyear, monthmonth, day1, hour0, minute0, second0, microsecond0, ) def bucket_key(value: datetime) - str: ##年份加月份作为索引依据 return f{value.year:04d}{value.month:02d} def table_name(value: datetime) - str: ##表名生成 return fandroid_messages_{bucket_key(value)} def create_message_table()##逻辑同最初单表逻辑这里省略 def register_bucket( ##注册新的表 cursor, key: str, name: str, start_at: datetime, end_at: datetime, ) - None: now utc_now() cursor.execute( INSERT INTO storage_registry ( bucket_key, table_name, start_at, end_at, status, schema_version, created_at, updated_at ) VALUES ( %s, %s, %s, %s, %s, %s, %s, %s ) ON DUPLICATE KEY UPDATE table_name VALUES(table_name), start_at VALUES(start_at), end_at VALUES(end_at), schema_version VALUES(schema_version), updated_at VALUES(updated_at) , ( key, name, start_at, end_at, active, SCHEMA_VERSION, now, now, ), ) def retire_old_buckets( ##旧表退出当前查询逻辑但仍能再后续再次查询 cursor, current_month: datetime, ) - None: cutoff add_months( current_month, -RETENTION_MONTHS, ) cursor.execute( UPDATE storage_registry SET status %s, updated_at %s WHERE end_at %s AND status %s , ( retired, utc_now(), cutoff, active, ), ) def ensure_registry(cursor) - None: cursor.execute(CREATE_REGISTRY_SQL) def manage_once() - None: ##执行函数将上方的各类子函数调用让manager融入当前通讯逻辑 connection None cursor None try: connection mysql.connector.connect( ##链接逻辑 **DB_CONFIG ) cursor connection.cursor() ##游标 ensure_registry(cursor) ##确认初始化合规 current_month month_start(utc_now()) ##声名最近月份 for offset in range( ##看似声名各类参数如表的生命周期 0, PRE_CREATE_MONTHS 1, ): start_at add_months( current_month, offset, ) end_at add_months( start_at, 1, ) key bucket_key(start_at) ##总表的索引生成 name table_name(start_at) ##新表的建立 create_message_table( ##复用最开始的单表中的message结构 cursor, name, ) register_bucket( ##注册 cursor, key, name, start_at, end_at, ) print( fensured bucket{key} ftable{name} ) retire_old_buckets( ##清理 cursor, current_month, ) connection.commit() except mysql.connector.Error as error: ##debug日志 if connection is not None: connection.rollback() print(fstorage manager error: {error}) raise finally: ##关闭链接停止对MySQL的操作 if cursor is not None: cursor.close() if connection is not None: connection.close() def main() - None: ##执行这里省略结合上方的代码一个简单的storage_manager就具备了基本的逻辑通过同步修改consumer新增对应的manager_router,我们就可以实现MySQL表的自更新获得一个可以不大规模改索引的基本条件——当然也为后续最大存储量maxsize_manager,不同用户分表user_manager提供了思路上的启发。manager_router这里不再给出但你可以这样想Linux-server中包含的consumer和storage_manager都是再同一个Linux共享内存中说以你大可以把它当作一个简单的函数调用和C编程中include头文件然后调用相应class中的函数很像二.一个简单的测试当代网络相关应用的很多技术其实都在围绕着一个问题如何处理业务场景下多个clientgateway和server之间信息的传递和不同组件中的业务信息加工。我们的项目也不例外为了高拟真的目标不论是测压还是不同业务的解耦多个组件实例之间的协作都是个不错的解决方案——回顾项目的基本通讯骨架也就是NATPort配合VM平台Linux的loaclIP雷电模拟器的LDPlayer多开模式reverse端口反向映射TCP链接和port监听项目的多实例测试其实并不算难工程量其实主要集中在大量数据的获取配置文件的优化以及细节调试上1.有效数据的生成不论是MySQL更深层次的索引优化还是信道可靠性的检测很多功能都依赖于data的量级在缺乏样本量的情况下我们其实很难做出有效的设计毕竟我们不能“抛开剂量谈毒性”——下面我会通过扩展Android-client的组件来给出相应的解决方案##资源引用 ...... data class LoadTestConfig( ##结合之前基本Android-client的envelope构造 val baseUrl: String, ##不同的是我们在这里还引入了packet的自定义大小和数量 val endpoint: String /message/batch, val producerId: String, val sessionId: String, val accessToken: String, val packetSize: Int, val totalCount: Int, val batchSize: Int, val intervalMillis: Long 0L, ) data class LoadTestSnapshot( ##测试结果 val attempted: Int, val succeeded: Int, val failed: Int, val bytesSent: Long, ) class LoadTestStats { private val attempted AtomicInteger() ##状态标识 private val succeeded AtomicInteger() private val failed AtomicInteger() private val bytesSent AtomicLong() fun addAttempt(count: Int, bytes: Long) { ##获取总状态 attempted.addAndGet(count) bytesSent.addAndGet(bytes) } fun addSuccess(count: Int) { ##成功数目 succeeded.addAndGet(count) } fun addFailure(count: Int) { ##失败数目 failed.addAndGet(count) } fun snapshot(): LoadTestSnapshot { ##相关结果传递 return LoadTestSnapshot( attempted attempted.get(), succeeded succeeded.get(), failed failed.get(), bytesSent bytesSent.get(), ) } } interface LoadTestListener { ##显示结果 fun onProgress(snapshot: LoadTestSnapshot) fun onFinished( snapshot: LoadTestSnapshot, error: Throwable?, ) } class LoadTestRunner( private val client: OkHttpClient, private val packetGenerator: PacketGenerator PacketGenerator(), ) { ##一些okhttp编程可以参考往期文章 private val jsonMediaType application/json; charsetutf-8.toMediaType() private val running AtomicBoolean(false) private val sequence AtomicInteger(0) Volatile private var activeCall: Call? null Volatile private var executor: ExecutorService? null fun start( config: LoadTestConfig, ##声名具体设置 listener: LoadTestListener, ) { require(config.packetSize in 1..PacketGenerator.MAX_PACKET_SIZE) ##发送条件 require(config.totalCount 0) require(config.batchSize 0) if (!running.compareAndSet(false, true)) { ##状态显示 throw IllegalStateException( load test is already running ) } sequence.set(0) ##占位初始化 val stats LoadTestStats() val worker Executors.newSingleThreadExecutor() executor worker worker.execute { var error: Throwable? null try { runLoop( ##参数声名 config config, stats stats, listener listener, ) } catch (throwable: Throwable) { ##状态日志 error throwable stats.addFailure( config.totalCount - stats.snapshot().succeeded ) } finally { ##详细日志 activeCall null running.set(false) listener.onFinished( stats.snapshot(), error, ) worker.shutdown() ##停止发送 executor null } } } private fun runLoop( config: LoadTestConfig, stats: LoadTestStats, listener: LoadTestListener, ) { var sent 0 while ( ##循环条件 running.get() sent config.totalCount ) { val currentBatchSize minOf( ##剩余待发送数目 config.batchSize, config.totalCount - sent, ) val batchBody buildBatchBody( ##datapacket构造 config config, startSequence sequence.get(), count currentBatchSize, ) sequence.addAndGet(currentBatchSize) ##占位更新 stats.addAttempt( ##状态日志 count currentBatchSize, bytes currentBatchSize.toLong() * config.packetSize, ) val request Request.Builder( .build() val call client.newCall(request) activeCall call try {} ##okhttp基本发送逻辑参考往期文章这里省略 stats.addSuccess(currentBatchSize) sent currentBatchSize ##状态更新 listener.onProgress( stats.snapshot() ) if ( ##缓冲时间 running.get() config.intervalMillis 0 ) { Thread.sleep( config.intervalMillis ) } } catch (error: IOException) { ##异常捕获 if (running.get()) { stats.addFailure(currentBatchSize) throw error } break } finally { if (activeCall call) { activeCall null } } } } private fun buildBatchBody( ##packet构造 config: LoadTestConfig, startSequence: Int, count: Int, ): JSONObject { val messages JSONArray() repeat(count) { offset - ##重复发送packet设置该次重复的参数 val message JSONObject() .put( message_id, UUID.randomUUID().toString(), ) .put( sequence, startSequence offset, ) .put( packet_size, config.packetSize, ) .put( payload, packetGenerator.generateAscii( config.packetSize ), ) messages.put(message) } return JSONObject() ##packet具体内容的获取 .put(type, batch) .put( run_id, config.producerId - config.sessionId, ) .put( batch_id, UUID.randomUUID().toString(), ) .put(producer_id, config.producerId) .put(session_id, config.sessionId) .put(access_token, config.accessToken) .put(message_count, count) .put(messages, messages) } fun stop() { ##停止 running.set(false) activeCall?.cancel() activeCall null executor?.shutdownNow() executor null } }当然这只是部分代码如果你想实现完整功能MainActivity应当同步增加入口同时XML布局控件也是必要的下面我会给出相应的packet包填充逻辑应该可以让你更好的了解packet的生成如果需要你也可以改变具体逻辑从而更加符合你个人需要的场景class PacketGenerator( ##packet生成器 private val random: Random Random(System.nanoTime()), ) { private val alphabet abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789 ##从上方源数据中随机选取并填充进packet fun generateAscii(sizeBytes: Int): String { require(sizeBytes in 1..MAX_PACKET_SIZE) { ##MAX_size防止过大消耗有限的性能 packet size must be between 1 and $MAX_PACKET_SIZE } val builder StringBuilder(sizeBytes) repeat(sizeBytes) { val index random.nextInt(alphabet.length) builder.append(alphabet[index]) } return builder.toString() } companion object { const val MAX_PACKET_SIZE 16 * 1024 } }2.配置优化之前我们一直使用JSON文件作为配置区在很多情况中这其实没什么问题但随着我们的业务变得更加复杂单个user往往不是单一持久权限而是由于session导致的在多个session中不同的权限对单个user——就好比你充VIP也能主动降低权和non-VIP用户一起用免费功能user的配置继续依靠JSON来做会变得比较臃肿而且加载JSON文件本身并不是很灵活——.env文件是个不错的替代品。通过配置不同版本的.env文件并妥善管理我们可以通过加载不同版本来实现不同配置的改变同时编写.env文件也不需要想编写JSON一样面对茫茫多的嵌套通过加载不同配置不同client的权限可以很快的变得同时结合token的动态计算session的权限也能跟着变化。.env文件的编写由于比较简单我不会给出相关的示例但注意使用这种方案需要我们先加载.env再启动相应模块而且需要注意加载顺序而且.env并不能够做动态配置相关工作还是由token机制以及本篇新加的storage_manager来做的对于具体的测试我建议你翻阅前一篇文章那里有比较完整的链路搭建虽然是单个实例但得益于LD的多开器和自带的各种工具多实例的运行并没有比单实例难无非是重复几遍取决于你想要的数量的事情这里我不多赘述。但是不要太多不然过多虚拟机可能会导致电脑性能受损2-5应该是个合理区间结合上文的批量packet获取有效数量级的data是比较容易的。尾记作为一个长期项目最难的部分已经开始不是技术本身了——如何在学业生活课外实践算法练习面试准备等一堆事情同时挤压得情况下平衡才是我的主要挑战——最近一直在刷算法题占了不少时间但现在能自己写写题解也算是个进步。现在是国庆假期只能说生活不易啊~为了未来能谋个生活还是太难了——明天奖励自己一个烂醉吧当然也祝看到这的你不用过上这种“享大福”的日子。下一步估计会尝试把一些算法运用到数据库上试一试或者模拟一些DNS处理逻辑。