已合并
【msServiceProfiler】【parse】【微重构】 增强detokenize事件的处理逻辑 #173
【msServiceProfiler】【parse】【微重构】 增强detokenize事件的处理逻辑 #173
已合并
Masaka创建于 3月6日
共 2 个文件变更+238-34
@@ -43,6 +43,9 @@ from ms_service_profiler.utils.trace_to_db import (
43 43 
44# 常量定义44# 常量定义
45LARGE_EVENTS_THRESHOLD = 5000000 # 当事件数量超过500万时,启用大数据量处理模式45LARGE_EVENTS_THRESHOLD = 5000000 # 当事件数量超过500万时,启用大数据量处理模式
46+DETOKENIZE_NAME = 'detokenize'
47+BATCH_FRAMEWORK_PROCESSING_NAME = 'batchFrameworkProcessing'
48+DETOKENIZE_FLOW_SUFFIX = '_detokenize'
46MAX_PROCESSES_LARGE = 8 # 大数据量处理时的最大进程数限制,避免过多进程导致系统资源耗尽49MAX_PROCESSES_LARGE = 8 # 大数据量处理时的最大进程数限制,避免过多进程导致系统资源耗尽
47MIN_CHUNK_SIZE_LARGE = 200000 # 大数据量处理时每个数据块的最小事件数量,确保每个进程有足够的工作量50MIN_CHUNK_SIZE_LARGE = 200000 # 大数据量处理时每个数据块的最小事件数量,确保每个进程有足够的工作量
48PROGRESS_REPORT_INTERVAL = 10 # 进度报告间隔,每处理10个数据块报告一次进度(仅限debug模式)51PROGRESS_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]
Masaka
MasakaMasaka3月8日

【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")
likedislike
408+ 
409+ 
369def add_flow_event(flow_event_df):410def 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=x427+ 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'] == 0430+ # 无 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).index436+ if main_events:
390- last_occurrences = exploded_df[multi_mask].groupby('rid').tail(1).index437+ 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]
Masaka
MasakaMasaka3月8日

【review】【性能】

使用列表推导式构建 detokenize_indices 的写法比较简洁,但在循环体内每次迭代都调用 e.get('name'),在事件数量较大的场景下会有轻微的重复调用开销。可参考改进代码:

enriched = [(i, e, e.get('name')) for i, e in enumerate(events_ordered)]
detokenize_indices = [i for i, e, name in enriched if name == DETOKENIZE_NAME]
likedislike
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_events465 return flow_trace_events
406 466 
407 467 
@@ -29,14 +29,19 @@ import pytest
29import pandas as pd29import pandas as pd
30 30 
31from ms_service_profiler.utils.file_open_check import OpenException31from 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_pid39+ 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)