已合并
【msServiceProfiler】【parse】【微重构】 增强detokenize事件的处理逻辑 #173
Masaka创建于 3月6日
【msServiceProfiler】【parse】【微重构】 增强detokenize事件的处理逻辑 #173
已合并
共 2 个文件变更+238-34
| @@ -43,6 +43,9 @@ from ms_service_profiler.utils.trace_to_db import ( | |||
| 43 | 43 | ||
| 44 | # 常量定义 | 44 | # 常量定义 |
| 45 | LARGE_EVENTS_THRESHOLD = 5000000 # 当事件数量超过500万时,启用大数据量处理模式 | 45 | LARGE_EVENTS_THRESHOLD = 5000000 # 当事件数量超过500万时,启用大数据量处理模式 |
| 46 | +DETOKENIZE_NAME = 'detokenize' | ||
| 47 | +BATCH_FRAMEWORK_PROCESSING_NAME = 'batchFrameworkProcessing' | ||
| 48 | +DETOKENIZE_FLOW_SUFFIX = '_detokenize' | ||
| 46 | MAX_PROCESSES_LARGE = 8 # 大数据量处理时的最大进程数限制,避免过多进程导致系统资源耗尽 | 49 | MAX_PROCESSES_LARGE = 8 # 大数据量处理时的最大进程数限制,避免过多进程导致系统资源耗尽 |
| 47 | MIN_CHUNK_SIZE_LARGE = 200000 # 大数据量处理时每个数据块的最小事件数量,确保每个进程有足够的工作量 | 50 | MIN_CHUNK_SIZE_LARGE = 200000 # 大数据量处理时每个数据块的最小事件数量,确保每个进程有足够的工作量 |
| 48 | PROGRESS_REPORT_INTERVAL = 10 # 进度报告间隔,每处理10个数据块报告一次进度(仅限debug模式) | 51 | PROGRESS_REPORT_INTERVAL = 10 # 进度报告间隔,每处理10个数据块报告一次进度(仅限debug模式) |
| @@ -366,42 +369,99 @@ def save_trace_data_into_json(trace_data, output): | |||
| 366 | write_trace_data_to_file(trace_data.get("traceEvents", []), file_path) | 369 | write_trace_data_to_file(trace_data.get("traceEvents", []), file_path) |
| 367 | 370 | ||
| 368 | 371 | ||
| 372 | +def _flow_events_from_sequence(events, flow_id, rid_for_cat, single_event_as_f=False): | ||
| 373 | + """将事件序列转为 flow 的 s/t/f 事件列表。events 为按 start_time 排序的行 dict 列表。 | ||
| 374 | + single_event_as_f:当序列仅一个事件时是否标为 'f'(用于 detokenize 分支终点)。""" | ||
| 375 | + result = [] | ||
| 376 | + for i, row in enumerate(events): | ||
| 377 | + ph = 't' | ||
| 378 | + if len(events) == 1: | ||
| 379 | + ph = 'f' if single_event_as_f else ('i' if row.get('during_time', 0) == 0 else 'X') | ||
| 380 | + else: | ||
| 381 | + if i == 0: | ||
| 382 | + ph = 's' | ||
| 383 | + elif i == len(events) - 1: | ||
| 384 | + ph = 'f' | ||
| 385 | + result.append({ | ||
| 386 | + 'name': 'flow_' + str(flow_id), | ||
| 387 | + 'ph': ph, | ||
| 388 | + 'ts': row['start_time'], | ||
| 389 | + 'id': str(flow_id), | ||
| 390 | + 'cat': str(rid_for_cat), | ||
| 391 | + 'pid': row['pid'], | ||
| 392 | + 'tid': row['domain'], | ||
| 393 | + }) | ||
| 394 | + return result | ||
| 395 | + | ||
| 396 | + | ||
| 397 | +def _find_detokenize_predecessor(events_before_detokenize): | ||
| 398 | + """ | ||
| 399 | + 在 detokenize 之前的事件中找分支前驱:优先最近的非 batchFrameworkProcessing,否则取时间上最近的一个。 | ||
| 400 | + events_before_detokenize 已按 start_time 升序。 | ||
| 401 | + """ | ||
| 402 | + if not events_before_detokenize: | ||
| 403 | + return None | ||
| 404 | + non_batch = [e for e in events_before_detokenize if e.get('name') != BATCH_FRAMEWORK_PROCESSING_NAME] | ||
| 405 | + if non_batch: | ||
| 406 | + return max(non_batch, key=lambda e: e['start_time']) | ||
| 407 | + return events_before_detokenize[-1] | ||
| 408 | + | ||
| 409 | + | ||
| 369 | def add_flow_event(flow_event_df): | 410 | def add_flow_event(flow_event_df): |
| 370 | - flow_event_df.loc[:, 'rid'] = flow_event_df['rid'].str.split(',') | 411 | + if flow_event_df is None or flow_event_df.empty: |
| 412 | + return [] | ||
| 413 | + flow_event_df = flow_event_df.copy() | ||
| 414 | + # 兼容 rid 为字符串(逗号分隔)或已为列表 | ||
| 415 | + rid_col = flow_event_df['rid'] | ||
| 416 | + if rid_col.dtype == object and len(rid_col) > 0 and isinstance(rid_col.iloc[0], str): | ||
| 417 | + flow_event_df.loc[:, 'rid'] = flow_event_df['rid'].astype(str).str.split(',') | ||
| 371 | exploded_df = flow_event_df.explode('rid') | 418 | exploded_df = flow_event_df.explode('rid') |
| 372 | exploded_df['tid'] = exploded_df['domain'] | 419 | exploded_df['tid'] = exploded_df['domain'] |
| 373 | 420 | ||
| 374 | - # 初始化 ph 列为默认值 | 421 | + flow_trace_events = [] |
| 375 | - exploded_df['ph'] = 't' | 422 | + for rid, group in exploded_df.groupby('rid', sort=False): |
| 423 | + # 按 start_time 排序 | ||
| 424 | + sorted_group = group.sort_values('start_time') | ||
| 425 | + events_ordered = sorted_group.to_dict('records') | ||
| 376 | 426 | ||
| 377 | - # 如果某个 rid 只有一行,根据during_time是否为0区分,为0就是ph=i,否则ph=x | 427 | + has_detokenize = any(e.get('name') == DETOKENIZE_NAME for e in events_ordered) |
| 378 | - single_occurrences = exploded_df['rid'].value_counts() | ||
| 379 | - single_rids = single_occurrences[single_occurrences == 1].index | ||
| 380 | - single_mask = exploded_df['rid'].isin(single_rids) | ||
| 381 | 428 | ||
| 382 | - # 根据 during_time 的值来设置 ph(优先处理) | 429 | + if not has_detokenize: |
| 383 | - during_time_zero_mask = exploded_df['during_time'] == 0 | 430 | + # 无 detokenize:主链一条,id=rid,与原有行为一致 |
| 384 | - exploded_df.loc[single_mask & during_time_zero_mask, 'ph'] = 'i' # 瞬时事件 | 431 | + flow_trace_events.extend(_flow_events_from_sequence(events_ordered, rid, rid)) |
| 385 | - exploded_df.loc[single_mask & ~during_time_zero_mask, 'ph'] = 'X' # 完整事件 | 432 | + continue |
| 386 | 433 | ||
| 387 | - # 找出每个 rid 的第一次和最后一次出现的位置(排除单次出现的) | 434 | + # 有 detokenize:主链排除 detokenize;每个 detokenize 事件各生成一条分支链(predecessor -> 该 detokenize) |
| 388 | - multi_mask = ~single_mask # 多次出现的记录 | 435 | + main_events = [e for e in events_ordered if e.get('name') != DETOKENIZE_NAME] |
| 389 | - first_occurrences = exploded_df[multi_mask].groupby('rid').head(1).index | 436 | + if main_events: |
| 390 | - last_occurrences = exploded_df[multi_mask].groupby('rid').tail(1).index | 437 | + flow_trace_events.extend(_flow_events_from_sequence(main_events, rid, rid)) |
| 391 | 438 | ||
| 392 | - # 设置第一次出现为 's' | 439 | + # 分支:对每个 detokenize 事件生成一条分支,保证每个 detokenize 都有异步连线;分支内不包含其他 detokenize,避免「detokenize 后又连到后续事件」 |
| 393 | - exploded_df.loc[first_occurrences, 'ph'] = 's' | 440 | + detokenize_indices = [i for i, e in enumerate(events_ordered) if e.get('name') == DETOKENIZE_NAME] |
【review】【性能】 使用列表推导式构建
![]() ![]() | |||
| 441 | + for det_ordinal, det_idx in enumerate(detokenize_indices): | ||
| 442 | + events_before = events_ordered[:det_idx] | ||
| 443 | + pred = _find_detokenize_predecessor(events_before) | ||
| 394 | 444 | ||
| 395 | - # 设置最后一次出现为 'f' | 445 | + if pred is None: |
| 396 | - exploded_df.loc[last_occurrences, 'ph'] = 'f' | 446 | + branch_events = [events_ordered[det_idx]] |
| 447 | + else: | ||
| 448 | + try: | ||
| 449 | + pred_idx = next( | ||
| 450 | + i for i, e in enumerate(events_ordered) | ||
| 451 | + if e.get('name') == pred.get('name') and e.get('start_time') == pred.get('start_time') | ||
| 452 | + ) | ||
| 453 | + except StopIteration: | ||
| 454 | + pred_idx = det_idx - 1 # pred 来自 events_before,理论上不会发生 | ||
| 455 | + # 只取起点到 pred 之间的非 detokenize 事件 + 当前 detokenize,避免把前面的 detokenize 带进分支导致图上「detokenize 连到后续事件」 | ||
| 456 | + segment_before_this_det = events_ordered[:pred_idx + 1] | ||
| 457 | + branch_events = [e for e in segment_before_this_det if e.get('name') != DETOKENIZE_NAME] + [events_ordered[det_idx]] | ||
| 458 | + | ||
| 459 | + # 多个 detokenize 时用序号区分 flow id,便于在 tracing 图中分别显示 | ||
| 460 | + branch_flow_id = str(rid) + DETOKENIZE_FLOW_SUFFIX + ('' if len(detokenize_indices) == 1 else '_' + str(det_ordinal)) | ||
| 461 | + flow_trace_events.extend(_flow_events_from_sequence( | ||
| 462 | + branch_events, branch_flow_id, rid, single_event_as_f=True | ||
| 463 | + )) | ||
| 397 | 464 | ||
| 398 | - exploded_df['bp'] = ['b' if ph == 'f' else '' for ph in exploded_df['ph']] | ||
| 399 | - exploded_df['name'] = 'flow_' + exploded_df['rid'] | ||
| 400 | - exploded_df['ts'] = exploded_df['start_time'] | ||
| 401 | - exploded_df['id'] = exploded_df['rid'] | ||
| 402 | - exploded_df['cat'] = exploded_df['rid'] | ||
| 403 | - exploded_df['pid'] = exploded_df['pid'] | ||
| 404 | - flow_trace_events = exploded_df[['name', 'ph', 'ts', 'id', 'cat', 'pid', 'tid']].to_dict(orient='records') | ||
| 405 | return flow_trace_events | 465 | return flow_trace_events |
| 406 | 466 | ||
| 407 | 467 | ||
| @@ -29,14 +29,19 @@ import pytest | |||
| 29 | import pandas as pd | 29 | import pandas as pd |
| 30 | 30 | ||
| 31 | from ms_service_profiler.utils.file_open_check import OpenException | 31 | from ms_service_profiler.utils.file_open_check import OpenException |
| 32 | -from ms_service_profiler.exporters.exporter_trace import ExporterTrace, write_trace_data_to_file, \ | 32 | +from ms_service_profiler.exporters.exporter_trace import ( |
| 33 | - save_trace_data_into_json, add_flow_event, create_trace_events, sort_trace_events_by_tid, add_mem_events, \ | 33 | + ExporterTrace, write_trace_data_to_file, |
| 34 | - load_single_prof, find_cann_pid, merge_json_data, add_npu_events, add_kvcache_events, add_cpu_events, \ | 34 | + save_trace_data_into_json, add_flow_event, create_trace_events, sort_trace_events_by_tid, add_mem_events, |
| 35 | - add_pull_kvcache_events, _prepare_data_smart_parallel, save_trace_data_into_db, _build_track_id_mapping_smart, \ | 35 | + DETOKENIZE_NAME, BATCH_FRAMEWORK_PROCESSING_NAME, DETOKENIZE_FLOW_SUFFIX, |
| 36 | - _setup_database_optimizations, _collect_pid_tid_and_meta_events, _get_tid_from_event, _process_meta_events_batch, \ | 36 | +) |
| 37 | - _find_thread_sort_index, _process_chunk_smart, _process_single_event, _is_slice_event, _is_counter_event, \ | 37 | +from ms_service_profiler.exporters.exporter_trace import ( |
| 38 | - _is_flow_event, _prepare_slice_data_smart, _prepare_counter_data_smart, _prepare_flow_data_smart, \ | 38 | + load_single_prof, find_cann_pid, merge_json_data, add_npu_events, add_kvcache_events, add_cpu_events, |
| 39 | - sort_trace_events_by_pid | 39 | + add_pull_kvcache_events, _prepare_data_smart_parallel, save_trace_data_into_db, _build_track_id_mapping_smart, |
| 40 | + _setup_database_optimizations, _collect_pid_tid_and_meta_events, _get_tid_from_event, _process_meta_events_batch, | ||
| 41 | + _find_thread_sort_index, _process_chunk_smart, _process_single_event, _is_slice_event, _is_counter_event, | ||
| 42 | + _is_flow_event, _prepare_slice_data_smart, _prepare_counter_data_smart, _prepare_flow_data_smart, | ||
| 43 | + sort_trace_events_by_pid, | ||
| 44 | +) | ||
| 40 | 45 | ||
| 41 | 46 | ||
| 42 | # Mock 数据 | 47 | # Mock 数据 |
| @@ -1195,3 +1200,142 @@ class TestSortTraceEventsByPid(unittest.TestCase): | |||
| 1195 | result = _prepare_data_smart_parallel(events, {}) | 1200 | result = _prepare_data_smart_parallel(events, {}) |
| 1196 | 1201 | ||
| 1197 | self.assertEqual(len(result['slice']), 1) | 1202 | self.assertEqual(len(result['slice']), 1) |
| 1203 | + | ||
| 1204 | + | ||
| 1205 | +# ====================================================================== | ||
| 1206 | +# add_flow_event 双 flow 逻辑单测与边界用例 | ||
| 1207 | +# ====================================================================== | ||
| 1208 | + | ||
| 1209 | + | ||
| 1210 | +def _make_flow_event_df(rows, rid=0, pid=1): | ||
| 1211 | + """构造 flow_event_df:每行 name, domain, rid, start_time, end_time, during_time, pid.""" | ||
| 1212 | + records = [] | ||
| 1213 | + for r in rows: | ||
| 1214 | + name = r['name'] | ||
| 1215 | + start = r['start_time'] | ||
| 1216 | + dur = r.get('during_time', 10) | ||
| 1217 | + records.append({ | ||
| 1218 | + 'name': name, | ||
| 1219 | + 'domain': r.get('domain', 'Api'), | ||
| 1220 | + 'rid': rid, | ||
| 1221 | + 'start_time': start, | ||
| 1222 | + 'end_time': start + dur, | ||
| 1223 | + 'during_time': dur, | ||
| 1224 | + 'pid': r.get('pid', pid), | ||
| 1225 | + }) | ||
| 1226 | + return pd.DataFrame(records) | ||
| 1227 | + | ||
| 1228 | + | ||
| 1229 | +def _flows_by_id(flow_events): | ||
| 1230 | + """按 id 分组,每组按 ts 排序。""" | ||
| 1231 | + by_id = {} | ||
| 1232 | + for e in flow_events: | ||
| 1233 | + fid = e['id'] | ||
| 1234 | + by_id.setdefault(fid, []).append(e) | ||
| 1235 | + for fid in by_id: | ||
| 1236 | + by_id[fid].sort(key=lambda x: (x['ts'], x.get('ph', ''))) | ||
| 1237 | + return by_id | ||
| 1238 | + | ||
| 1239 | + | ||
| 1240 | +class TestAddFlowEventDualFlow(unittest.TestCase): | ||
| 1241 | + """add_flow_event 双 flow 行为:无 detokenize 仅主链;有 detokenize 时主链+分支链.""" | ||
| 1242 | + | ||
| 1243 | + def test_no_detokenize_main_chain_only(self): | ||
| 1244 | + """无 detokenize(仅主链):只生成一条 flow,id=rid,事件顺序与时间序一致,s/t/f 标法同原先。""" | ||
| 1245 | + rid = 42 | ||
| 1246 | + flow_event_df = _make_flow_event_df([ | ||
| 1247 | + {'name': 'modelExec', 'start_time': 100}, | ||
| 1248 | + {'name': BATCH_FRAMEWORK_PROCESSING_NAME, 'start_time': 200}, | ||
| 1249 | + ], rid=rid) | ||
| 1250 | + result = add_flow_event(flow_event_df) | ||
| 1251 | + flows = _flows_by_id(result) | ||
| 1252 | + self.assertEqual(len(flows), 1, '应只有一条 flow') | ||
| 1253 | + self.assertIn(str(rid), flows) | ||
| 1254 | + main = flows[str(rid)] | ||
| 1255 | + self.assertEqual(len(main), 2) | ||
| 1256 | + self.assertEqual(main[0]['ph'], 's') | ||
| 1257 | + self.assertEqual(main[1]['ph'], 'f') | ||
| 1258 | + self.assertEqual(main[0]['ts'], 100) | ||
| 1259 | + self.assertEqual(main[1]['ts'], 200) | ||
| 1260 | + | ||
| 1261 | + def test_main_chain_and_branch_model_exec_batch_detokenize(self): | ||
| 1262 | + """modelExec → batchFrameworkProcessing → detokenize:主链无 detokenize,分支 modelExec→detokenize,分支 id=rid_detokenize。""" | ||
| 1263 | + rid = 1 | ||
| 1264 | + flow_event_df = _make_flow_event_df([ | ||
| 1265 | + {'name': 'modelExec', 'start_time': 100}, | ||
| 1266 | + {'name': BATCH_FRAMEWORK_PROCESSING_NAME, 'start_time': 200}, | ||
| 1267 | + {'name': DETOKENIZE_NAME, 'start_time': 250}, | ||
| 1268 | + ], rid=rid) | ||
| 1269 | + result = add_flow_event(flow_event_df) | ||
| 1270 | + flows = _flows_by_id(result) | ||
| 1271 | + self.assertEqual(len(flows), 2, '应有主链 + 分支链') | ||
| 1272 | + self.assertIn(str(rid), flows) | ||
| 1273 | + self.assertIn(str(rid) + DETOKENIZE_FLOW_SUFFIX, flows) | ||
| 1274 | + main = flows[str(rid)] | ||
| 1275 | + branch = flows[str(rid) + DETOKENIZE_FLOW_SUFFIX] | ||
| 1276 | + self.assertEqual(len(main), 2, '主链应为 modelExec, batchFrameworkProcessing') | ||
| 1277 | + self.assertEqual(main[0]['ts'], 100) | ||
| 1278 | + self.assertEqual(main[1]['ts'], 200) | ||
| 1279 | + self.assertEqual(main[0]['ph'], 's') | ||
| 1280 | + self.assertEqual(main[1]['ph'], 'f') | ||
| 1281 | + self.assertEqual(len(branch), 2, '分支应为 modelExec, detokenize') | ||
| 1282 | + self.assertEqual(branch[0]['ts'], 100) | ||
| 1283 | + self.assertEqual(branch[1]['ts'], 250) | ||
| 1284 | + self.assertEqual(branch[0]['ph'], 's') | ||
| 1285 | + self.assertEqual(branch[1]['ph'], 'f') | ||
| 1286 | + | ||
| 1287 | + def test_main_chain_skips_detokenize_branch_model_exec_detokenize(self): | ||
| 1288 | + """modelExec → detokenize → batchFrameworkProcessing:主链跳过 detokenize,分支 modelExec→detokenize。""" | ||
| 1289 | + rid = 2 | ||
| 1290 | + flow_event_df = _make_flow_event_df([ | ||
| 1291 | + {'name': 'modelExec', 'start_time': 100}, | ||
| 1292 | + {'name': DETOKENIZE_NAME, 'start_time': 150}, | ||
| 1293 | + {'name': BATCH_FRAMEWORK_PROCESSING_NAME, 'start_time': 200}, | ||
| 1294 | + ], rid=rid) | ||
| 1295 | + result = add_flow_event(flow_event_df) | ||
| 1296 | + flows = _flows_by_id(result) | ||
| 1297 | + self.assertEqual(len(flows), 2) | ||
| 1298 | + main = flows[str(rid)] | ||
| 1299 | + branch = flows[str(rid) + DETOKENIZE_FLOW_SUFFIX] | ||
| 1300 | + self.assertEqual(len(main), 2, '主链应为 modelExec, batchFrameworkProcessing(时间序跳过 detokenize)') | ||
| 1301 | + self.assertEqual(main[0]['ts'], 100) | ||
| 1302 | + self.assertEqual(main[1]['ts'], 200) | ||
| 1303 | + self.assertEqual(len(branch), 2, '分支应为 modelExec, detokenize') | ||
| 1304 | + self.assertEqual(branch[0]['ts'], 100) | ||
| 1305 | + self.assertEqual(branch[1]['ts'], 150) | ||
| 1306 | + self.assertEqual(branch[1]['ph'], 'f') | ||
| 1307 | + | ||
| 1308 | + def test_degenerate_batch_only_before_detokenize(self): | ||
| 1309 | + """仅 batchFrameworkProcessing → detokenize(退化):分支链为 batch→detokenize。""" | ||
| 1310 | + rid = 3 | ||
| 1311 | + flow_event_df = _make_flow_event_df([ | ||
| 1312 | + {'name': BATCH_FRAMEWORK_PROCESSING_NAME, 'start_time': 100}, | ||
| 1313 | + {'name': DETOKENIZE_NAME, 'start_time': 200}, | ||
| 1314 | + ], rid=rid) | ||
| 1315 | + result = add_flow_event(flow_event_df) | ||
| 1316 | + flows = _flows_by_id(result) | ||
| 1317 | + self.assertEqual(len(flows), 2) | ||
| 1318 | + main = flows[str(rid)] | ||
| 1319 | + branch = flows[str(rid) + DETOKENIZE_FLOW_SUFFIX] | ||
| 1320 | + self.assertEqual(len(main), 1, '主链只有 batchFrameworkProcessing') | ||
| 1321 | + self.assertEqual(main[0]['ts'], 100) | ||
| 1322 | + self.assertEqual(len(branch), 2, '分支为 batch→detokenize') | ||
| 1323 | + self.assertEqual(branch[0]['ts'], 100) | ||
| 1324 | + self.assertEqual(branch[1]['ts'], 200) | ||
| 1325 | + self.assertEqual(branch[1]['ph'], 'f') | ||
| 1326 | + | ||
| 1327 | + def test_branch_only_detokenize_single_event_ph_f(self): | ||
| 1328 | + """分支仅含 detokenize(无前驱)时,该唯一事件 ph='f'。""" | ||
| 1329 | + rid = 4 | ||
| 1330 | + flow_event_df = _make_flow_event_df([ | ||
| 1331 | + {'name': DETOKENIZE_NAME, 'start_time': 100}, | ||
| 1332 | + ], rid=rid) | ||
| 1333 | + result = add_flow_event(flow_event_df) | ||
| 1334 | + flows = _flows_by_id(result) | ||
| 1335 | + self.assertEqual(len(flows), 1, '主链空不生成事件,仅有一条分支链') | ||
| 1336 | + self.assertNotIn(str(rid), flows, '主链无事件时不应产出主链 flow') | ||
| 1337 | + self.assertIn(str(rid) + DETOKENIZE_FLOW_SUFFIX, flows) | ||
| 1338 | + branch = flows[str(rid) + DETOKENIZE_FLOW_SUFFIX] | ||
| 1339 | + self.assertEqual(len(branch), 1) | ||
| 1340 | + self.assertEqual(branch[0]['ph'], 'f') | ||
| 1341 | + self.assertEqual(branch[0]['ts'], 100) | ||


【review】【逻辑】
这段代码对 detokenize 事件的处理逻辑做了重构,通过
_find_detokenize_predecessor函数优先选择非 batchFrameworkProcessing 的事件作为前驱,设计思路是清晰的。不过在第 407 行 fallback 取时间最近事件作为前驱的这个分支,根据业务逻辑应该极少被触发,建议增加一条日志记录来便于后续排查确认该分支是否被触发。可参考改进代码:# 在 return events_before_detokenize[-1] 之前加 logging.debug("No non-batchFrameworkProcessing event found, using last event as predecessor")