perf: 18个v1 DSPY全改用llmid缓存,去3表JOIN+DB连接

This commit is contained in:
yumoqing 2026-07-24 11:01:40 +08:00
parent ef1e228f19
commit ad5d8c00c5
19 changed files with 612 additions and 562 deletions

163
load_test.py Normal file
View File

@ -0,0 +1,163 @@
#!/usr/bin/env python3
"""并发压力测试 llmage /v1/chat/completions统计 TTFB / 完成时间 / QPM"""
import asyncio
import aiohttp
import time
import json
import sys
import statistics
from dataclasses import dataclass, field
from typing import List
URL = "https://token.opencomputing.cn/llmage/v1/chat/completions"
TOKEN = "V9J41PngWBUU6gdHWJWDJ"
MODEL = "qwen3.6-35b-a3b"
DURATION = 180 # 3 分钟
CONCURRENCIES = [10, 50, 100, 200]
@dataclass
class ReqStat:
prompt_idx: int
start_ts: float
first_byte_ts: float | None = None
end_ts: float | None = None
async def worker(session: aiohttp.ClientSession, idx: int, stats_out: list):
"""单个请求:发送 stream 请求,记录首字时间和完成时间"""
prompt = f"请用一句话介绍你自己,编号{idx}"
payload = {
"model": MODEL,
"stream": True,
"messages": [{"role": "user", "content": prompt}],
}
stat = ReqStat(prompt_idx=idx, start_ts=time.monotonic())
try:
async with session.post(
URL,
json=payload,
headers={
"Content-Type": "application/json",
"Authorization": f"Bearer {TOKEN}",
},
timeout=aiohttp.ClientTimeout(total=120),
) as resp:
first = True
async for line in resp.content:
if first:
stat.first_byte_ts = time.monotonic()
first = False
# 读完所有 chunk 才算完成
stat.end_ts = time.monotonic()
except Exception as e:
# 异常请求也记录TTFB=None 表示失败)
stat.end_ts = time.monotonic()
stats_out.append(stat)
async def run_concurrency(concurrency: int):
"""以固定并发运行 DURATION 秒,持续发起新请求"""
stats: List[ReqStat] = []
idx = 0
stop_at = time.monotonic() + DURATION
connector = aiohttp.TCPConnector(limit=concurrency + 20, force_close=True)
async with aiohttp.ClientSession(connector=connector) as session:
tasks: list[asyncio.Task] = []
while time.monotonic() < stop_at:
# 保持并发数:补满到 concurrency
while len(tasks) < concurrency and time.monotonic() < stop_at:
idx += 1
tasks.append(
asyncio.create_task(worker(session, idx, stats))
)
if not tasks:
break
# 等待任意一个完成,腾出槽位
done, tasks = await asyncio.wait(
tasks, return_when=asyncio.FIRST_COMPLETED, timeout=0.5
)
# 清理已完成的
tasks = list(tasks)
# 时间到,等待所有进行中的请求完成
if tasks:
await asyncio.wait(tasks)
return stats
def analyze(name: str, stats: List[ReqStat]):
"""分析并打印统计"""
ttfb_list = [s.first_byte_ts - s.start_ts for s in stats if s.first_byte_ts]
total_list = [s.end_ts - s.start_ts for s in stats if s.end_ts and s.first_byte_ts]
failed = sum(1 for s in stats if s.first_byte_ts is None)
total_req = len(stats)
elapsed = DURATION
qpm = total_req / (elapsed / 60)
print(f"\n{'='*60}")
print(f" 并发={name} | 运行{DURATION}s | 总请求={total_req} | 失败={failed}")
print(f"{'='*60}")
if ttfb_list:
print(f" TTFB (s): min={min(ttfb_list):.3f} avg={statistics.mean(ttfb_list):.3f} "
f"p50={statistics.median(ttfb_list):.3f} p95={_pct(ttfb_list, 95):.3f} p99={_pct(ttfb_list, 99):.3f}")
if total_list:
print(f" 完成 (s): min={min(total_list):.3f} avg={statistics.mean(total_list):.3f} "
f"p50={statistics.median(total_list):.3f} p95={_pct(total_list, 95):.3f} p99={_pct(total_list, 99):.3f}")
print(f" QPM: {qpm:.1f}")
print(f" QPS: {total_req / elapsed:.1f}")
# 按分钟分段统计
for minute in range(int(elapsed / 60)):
win_start = minute * 60
win_end = (minute + 1) * 60
cnt = sum(1 for s in stats if s.end_ts and (s.end_ts - s.start_ts) >= 0
and win_start <= (s.start_ts - stats[0].start_ts) < win_end)
print(f"{minute+1}分钟完成请求数: {cnt}")
return {
"concurrency": name,
"total": total_req,
"failed": failed,
"ttfb_avg": statistics.mean(ttfb_list) if ttfb_list else None,
"ttfb_p50": statistics.median(ttfb_list) if ttfb_list else None,
"ttfb_p95": _pct(ttfb_list, 95) if ttfb_list else None,
"total_avg": statistics.mean(total_list) if total_list else None,
"total_p50": statistics.median(total_list) if total_list else None,
"qpm": qpm,
}
def _pct(data, p):
return sorted(data)[int(len(data) * p / 100)]
async def main():
results = []
for c in CONCURRENCIES:
print(f"\n>>> 开始测试 并发={c} ...")
stats = await run_concurrency(c)
r = analyze(str(c), stats)
results.append(r)
# 汇总表格
print(f"\n{'='*60}")
print(" 汇总对比")
print(f"{'='*60}")
print(f" {'并发':>6} {'总请求':>8} {'失败':>5} {'TTFB_avg':>9} {'TTFB_p50':>9} {'TTFB_p95':>9} {'完成_avg':>9} {'QPM':>8}")
for r in results:
print(f" {r['concurrency']:>6} {r['total']:>8} {r['failed']:>5} "
f"{r['ttfb_avg']:.3f}s" if r['ttfb_avg'] else "N/A".rjust(9) + " "
f"{(r['ttfb_p50'] or 0):.3f}s".rjust(9) + " "
f"{(r['ttfb_p95'] or 0):.3f}s".rjust(9) + " "
f"{r['total_avg']:.3f}s".rjust(9) if r['total_avg'] else "N/A".rjust(9))
if __name__ == "__main__":
asyncio.run(main())

View File

@ -41,22 +41,11 @@ if not params_kw.prompt:
lctype = params_kw.catelogid lctype = params_kw.catelogid
env = request._run_ns env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: llmid = await env.get_llmid_cached(env, params_kw.model, lctype)
# Look up llm by model name and catalog type through llm_api_map if not llmid:
sql = """select distinct a.* from llm a debug(f'{params_kw.model=} not found for catalog {lctype}')
join llm_api_map m on a.id = m.llmid return openai_400()
join llmcatelog b on m.llmcatelogid = b.id params_kw.llmid = llmid
where (b.id = ${lctype}$ OR b.name = ${lctype}$)
and a.name=${model}$
and a.status = 'published'"""
recs = await sor.sqlExe(sql, {
'lctype': lctype,
'model': params_kw.model
})
if len(recs) == 0:
debug(f'{params_kw.model=} not found for catalog {lctype}')
return openai_400()
params_kw.llmid = recs[0].id
params_kw.llmcatelogid = lctype params_kw.llmcatelogid = lctype
debug(f'{params_kw.llmid=}') debug(f'{params_kw.llmid=}')

View File

@ -38,22 +38,11 @@ if not params_kw.audio_file:
lctype = params_kw.catelogid lctype = params_kw.catelogid
env = request._run_ns env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: llmid = await env.get_llmid_cached(env, params_kw.model, lctype)
# Look up llm by model name and catalog type through llm_api_map if not llmid:
sql = """select distinct a.* from llm a debug(f'{params_kw.model=} not found for catalog {lctype}')
join llm_api_map m on a.id = m.llmid return openai_400()
join llmcatelog b on m.llmcatelogid = b.id params_kw.llmid = llmid
where (b.id = ${lctype}$ OR b.name = ${lctype}$)
and a.name=${model}$
and a.status = 'published'"""
recs = await sor.sqlExe(sql, {
'lctype': lctype,
'model': params_kw.model
})
if len(recs) == 0:
debug(f'{params_kw.model=} not found for catalog {lctype}')
return openai_400()
params_kw.llmid = recs[0].id
params_kw.llmcatelogid = lctype params_kw.llmcatelogid = lctype
debug(f'{params_kw.llmid=}') debug(f'{params_kw.llmid=}')

View File

@ -39,22 +39,11 @@ if not params_kw.prompt:
lctype = params_kw.catelogid lctype = params_kw.catelogid
env = request._run_ns env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: llmid = await env.get_llmid_cached(env, params_kw.model, lctype)
# Look up llm by model name and catalog type through llm_api_map if not llmid:
sql = """select distinct a.* from llm a debug(f'{params_kw.model=} not found for catalog {lctype}')
join llm_api_map m on a.id = m.llmid return openai_400()
join llmcatelog b on m.llmcatelogid = b.id params_kw.llmid = llmid
where (b.id = ${lctype}$ OR b.name = ${lctype}$)
and a.name=${model}$
and a.status = 'published'"""
recs = await sor.sqlExe(sql, {
'lctype': lctype,
'model': params_kw.model
})
if len(recs) == 0:
debug(f'{params_kw.model=} not found for catalog {lctype}')
return openai_400()
params_kw.llmid = recs[0].id
params_kw.llmcatelogid = lctype params_kw.llmcatelogid = lctype
debug(f'{params_kw.llmid=}') debug(f'{params_kw.llmid=}')

View File

@ -1,34 +1,30 @@
# KTV Media API - asr-transcribe 1|# KTV Media API - asr-transcribe
# /v1/media/asr-transcribe 2|# /v1/media/asr-transcribe
3|
debug_params('params_kw', params_kw) 4|debug_params('params_kw', params_kw)
5|
userid = await get_user() 6|userid = await get_user()
userorgid = await get_userorgid() 7|userorgid = await get_userorgid()
if userid is None: 8|if userid is None:
return openai_403() 9| return openai_403()
10|
catelogid = 'ktv_pipeline' 11|catelogid = 'ktv_pipeline'
model_name = 'asr-transcribe' 12|model_name = 'asr-transcribe'
params_kw.model = model_name 13|params_kw.model = model_name
params_kw.catelogid = catelogid 14|params_kw.catelogid = catelogid
15|
env = request._run_ns 16|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 17|# llmid from cache (model+catelogid -> llmid)
sql = """select distinct a.* from llm a llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
join llm_api_map m on a.id = m.llmid if not llmid:
join llmcatelog b on m.llmcatelogid = b.id debug(f'model not found: params_kw.model')
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) return openai_400()
and a.model=${model}$ params_kw.llmid = llmid
and a.status = 'published'""" 28| params_kw.llmcatelogid = catelogid
recs = await sor.sqlExe(sql, {'catelogid': catelogid, 'model': model_name}) 29|
if len(recs) == 0: 30|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
return openai_400() 31|if not f:
params_kw.llmid = recs[0].id 32| return openai_429()
params_kw.llmcatelogid = catelogid 33|
34|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 35|
if not f:
return openai_429()
return await inference(request, env=env)

View File

@ -1,34 +1,30 @@
# KTV Media API - demucs-separate 1|# KTV Media API - demucs-separate
# /v1/media/demucs-separate 2|# /v1/media/demucs-separate
3|
debug_params('params_kw', params_kw) 4|debug_params('params_kw', params_kw)
5|
userid = await get_user() 6|userid = await get_user()
userorgid = await get_userorgid() 7|userorgid = await get_userorgid()
if userid is None: 8|if userid is None:
return openai_403() 9| return openai_403()
10|
catelogid = 'ktv_pipeline' 11|catelogid = 'ktv_pipeline'
model_name = 'demucs-separate' 12|model_name = 'demucs-separate'
params_kw.model = model_name 13|params_kw.model = model_name
params_kw.catelogid = catelogid 14|params_kw.catelogid = catelogid
15|
env = request._run_ns 16|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 17|# llmid from cache (model+catelogid -> llmid)
sql = """select distinct a.* from llm a llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
join llm_api_map m on a.id = m.llmid if not llmid:
join llmcatelog b on m.llmcatelogid = b.id debug(f'model not found: params_kw.model')
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) return openai_400()
and a.model=${model}$ params_kw.llmid = llmid
and a.status = 'published'""" 28| params_kw.llmcatelogid = catelogid
recs = await sor.sqlExe(sql, {'catelogid': catelogid, 'model': model_name}) 29|
if len(recs) == 0: 30|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
return openai_400() 31|if not f:
params_kw.llmid = recs[0].id 32| return openai_429()
params_kw.llmcatelogid = catelogid 33|
34|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 35|
if not f:
return openai_429()
return await inference(request, env=env)

View File

@ -1,34 +1,30 @@
# KTV Media API - face-compare 1|# KTV Media API - face-compare
# /v1/media/face-compare 2|# /v1/media/face-compare
3|
debug_params('params_kw', params_kw) 4|debug_params('params_kw', params_kw)
5|
userid = await get_user() 6|userid = await get_user()
userorgid = await get_userorgid() 7|userorgid = await get_userorgid()
if userid is None: 8|if userid is None:
return openai_403() 9| return openai_403()
10|
catelogid = 'ktv_pipeline' 11|catelogid = 'ktv_pipeline'
model_name = 'face-compare' 12|model_name = 'face-compare'
params_kw.model = model_name 13|params_kw.model = model_name
params_kw.catelogid = catelogid 14|params_kw.catelogid = catelogid
15|
env = request._run_ns 16|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 17|# llmid from cache (model+catelogid -> llmid)
sql = """select distinct a.* from llm a llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
join llm_api_map m on a.id = m.llmid if not llmid:
join llmcatelog b on m.llmcatelogid = b.id debug(f'model not found: params_kw.model')
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) return openai_400()
and a.model=${model}$ params_kw.llmid = llmid
and a.status = 'published'""" 28| params_kw.llmcatelogid = catelogid
recs = await sor.sqlExe(sql, {'catelogid': catelogid, 'model': model_name}) 29|
if len(recs) == 0: 30|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
return openai_400() 31|if not f:
params_kw.llmid = recs[0].id 32| return openai_429()
params_kw.llmcatelogid = catelogid 33|
34|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 35|
if not f:
return openai_429()
return await inference(request, env=env)

View File

@ -1,34 +1,30 @@
# KTV Media API - face-detect 1|# KTV Media API - face-detect
# /v1/media/face-detect 2|# /v1/media/face-detect
3|
debug_params('params_kw', params_kw) 4|debug_params('params_kw', params_kw)
5|
userid = await get_user() 6|userid = await get_user()
userorgid = await get_userorgid() 7|userorgid = await get_userorgid()
if userid is None: 8|if userid is None:
return openai_403() 9| return openai_403()
10|
catelogid = 'ktv_pipeline' 11|catelogid = 'ktv_pipeline'
model_name = 'face-detect' 12|model_name = 'face-detect'
params_kw.model = model_name 13|params_kw.model = model_name
params_kw.catelogid = catelogid 14|params_kw.catelogid = catelogid
15|
env = request._run_ns 16|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 17|# llmid from cache (model+catelogid -> llmid)
sql = """select distinct a.* from llm a llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
join llm_api_map m on a.id = m.llmid if not llmid:
join llmcatelog b on m.llmcatelogid = b.id debug(f'model not found: params_kw.model')
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) return openai_400()
and a.model=${model}$ params_kw.llmid = llmid
and a.status = 'published'""" 28| params_kw.llmcatelogid = catelogid
recs = await sor.sqlExe(sql, {'catelogid': catelogid, 'model': model_name}) 29|
if len(recs) == 0: 30|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
return openai_400() 31|if not f:
params_kw.llmid = recs[0].id 32| return openai_429()
params_kw.llmcatelogid = catelogid 33|
34|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 35|
if not f:
return openai_429()
return await inference(request, env=env)

View File

@ -1,34 +1,30 @@
# KTV Media API - face-recognize 1|# KTV Media API - face-recognize
# /v1/media/face-recognize 2|# /v1/media/face-recognize
3|
debug_params('params_kw', params_kw) 4|debug_params('params_kw', params_kw)
5|
userid = await get_user() 6|userid = await get_user()
userorgid = await get_userorgid() 7|userorgid = await get_userorgid()
if userid is None: 8|if userid is None:
return openai_403() 9| return openai_403()
10|
catelogid = 'ktv_pipeline' 11|catelogid = 'ktv_pipeline'
model_name = 'face-recognize' 12|model_name = 'face-recognize'
params_kw.model = model_name 13|params_kw.model = model_name
params_kw.catelogid = catelogid 14|params_kw.catelogid = catelogid
15|
env = request._run_ns 16|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 17|# llmid from cache (model+catelogid -> llmid)
sql = """select distinct a.* from llm a llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
join llm_api_map m on a.id = m.llmid if not llmid:
join llmcatelog b on m.llmcatelogid = b.id debug(f'model not found: params_kw.model')
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) return openai_400()
and a.model=${model}$ params_kw.llmid = llmid
and a.status = 'published'""" 28| params_kw.llmcatelogid = catelogid
recs = await sor.sqlExe(sql, {'catelogid': catelogid, 'model': model_name}) 29|
if len(recs) == 0: 30|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
return openai_400() 31|if not f:
params_kw.llmid = recs[0].id 32| return openai_429()
params_kw.llmcatelogid = catelogid 33|
34|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 35|
if not f:
return openai_429()
return await inference(request, env=env)

View File

@ -1,34 +1,30 @@
# KTV Media API - merge-video 1|# KTV Media API - merge-video
# /v1/media/merge-video 2|# /v1/media/merge-video
3|
debug_params('params_kw', params_kw) 4|debug_params('params_kw', params_kw)
5|
userid = await get_user() 6|userid = await get_user()
userorgid = await get_userorgid() 7|userorgid = await get_userorgid()
if userid is None: 8|if userid is None:
return openai_403() 9| return openai_403()
10|
catelogid = 'ktv_pipeline' 11|catelogid = 'ktv_pipeline'
model_name = 'merge-video' 12|model_name = 'merge-video'
params_kw.model = model_name 13|params_kw.model = model_name
params_kw.catelogid = catelogid 14|params_kw.catelogid = catelogid
15|
env = request._run_ns 16|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 17|# llmid from cache (model+catelogid -> llmid)
sql = """select distinct a.* from llm a llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
join llm_api_map m on a.id = m.llmid if not llmid:
join llmcatelog b on m.llmcatelogid = b.id debug(f'model not found: params_kw.model')
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) return openai_400()
and a.model=${model}$ params_kw.llmid = llmid
and a.status = 'published'""" 28| params_kw.llmcatelogid = catelogid
recs = await sor.sqlExe(sql, {'catelogid': catelogid, 'model': model_name}) 29|
if len(recs) == 0: 30|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
return openai_400() 31|if not f:
params_kw.llmid = recs[0].id 32| return openai_429()
params_kw.llmcatelogid = catelogid 33|
34|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 35|
if not f:
return openai_429()
return await inference(request, env=env)

View File

@ -1,34 +1,30 @@
# KTV Media API - realesrgan-upscale 1|# KTV Media API - realesrgan-upscale
# /v1/media/realesrgan-upscale 2|# /v1/media/realesrgan-upscale
3|
debug_params('params_kw', params_kw) 4|debug_params('params_kw', params_kw)
5|
userid = await get_user() 6|userid = await get_user()
userorgid = await get_userorgid() 7|userorgid = await get_userorgid()
if userid is None: 8|if userid is None:
return openai_403() 9| return openai_403()
10|
catelogid = 'ktv_pipeline' 11|catelogid = 'ktv_pipeline'
model_name = 'realesrgan-upscale' 12|model_name = 'realesrgan-upscale'
params_kw.model = model_name 13|params_kw.model = model_name
params_kw.catelogid = catelogid 14|params_kw.catelogid = catelogid
15|
env = request._run_ns 16|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 17|# llmid from cache (model+catelogid -> llmid)
sql = """select distinct a.* from llm a llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
join llm_api_map m on a.id = m.llmid if not llmid:
join llmcatelog b on m.llmcatelogid = b.id debug(f'model not found: params_kw.model')
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) return openai_400()
and a.model=${model}$ params_kw.llmid = llmid
and a.status = 'published'""" 28| params_kw.llmcatelogid = catelogid
recs = await sor.sqlExe(sql, {'catelogid': catelogid, 'model': model_name}) 29|
if len(recs) == 0: 30|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
return openai_400() 31|if not f:
params_kw.llmid = recs[0].id 32| return openai_429()
params_kw.llmcatelogid = catelogid 33|
34|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 35|
if not f:
return openai_429()
return await inference(request, env=env)

View File

@ -1,34 +1,30 @@
# KTV Media API - rvc-convert 1|# KTV Media API - rvc-convert
# /v1/media/rvc-convert 2|# /v1/media/rvc-convert
3|
debug_params('params_kw', params_kw) 4|debug_params('params_kw', params_kw)
5|
userid = await get_user() 6|userid = await get_user()
userorgid = await get_userorgid() 7|userorgid = await get_userorgid()
if userid is None: 8|if userid is None:
return openai_403() 9| return openai_403()
10|
catelogid = 'ktv_pipeline' 11|catelogid = 'ktv_pipeline'
model_name = 'rvc-convert' 12|model_name = 'rvc-convert'
params_kw.model = model_name 13|params_kw.model = model_name
params_kw.catelogid = catelogid 14|params_kw.catelogid = catelogid
15|
env = request._run_ns 16|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 17|# llmid from cache (model+catelogid -> llmid)
sql = """select distinct a.* from llm a llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
join llm_api_map m on a.id = m.llmid if not llmid:
join llmcatelog b on m.llmcatelogid = b.id debug(f'model not found: params_kw.model')
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) return openai_400()
and a.model=${model}$ params_kw.llmid = llmid
and a.status = 'published'""" 28| params_kw.llmcatelogid = catelogid
recs = await sor.sqlExe(sql, {'catelogid': catelogid, 'model': model_name}) 29|
if len(recs) == 0: 30|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
return openai_400() 31|if not f:
params_kw.llmid = recs[0].id 32| return openai_429()
params_kw.llmcatelogid = catelogid 33|
34|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 35|
if not f:
return openai_429()
return await inference(request, env=env)

View File

@ -1,34 +1,30 @@
# KTV Media API - songrate-evaluate 1|# KTV Media API - songrate-evaluate
# /v1/media/songrate-evaluate 2|# /v1/media/songrate-evaluate
3|
debug_params('params_kw', params_kw) 4|debug_params('params_kw', params_kw)
5|
userid = await get_user() 6|userid = await get_user()
userorgid = await get_userorgid() 7|userorgid = await get_userorgid()
if userid is None: 8|if userid is None:
return openai_403() 9| return openai_403()
10|
catelogid = 'ktv_pipeline' 11|catelogid = 'ktv_pipeline'
model_name = 'songrate-evaluate' 12|model_name = 'songrate-evaluate'
params_kw.model = model_name 13|params_kw.model = model_name
params_kw.catelogid = catelogid 14|params_kw.catelogid = catelogid
15|
env = request._run_ns 16|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 17|# llmid from cache (model+catelogid -> llmid)
sql = """select distinct a.* from llm a llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
join llm_api_map m on a.id = m.llmid if not llmid:
join llmcatelog b on m.llmcatelogid = b.id debug(f'model not found: params_kw.model')
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) return openai_400()
and a.model=${model}$ params_kw.llmid = llmid
and a.status = 'published'""" 28| params_kw.llmcatelogid = catelogid
recs = await sor.sqlExe(sql, {'catelogid': catelogid, 'model': model_name}) 29|
if len(recs) == 0: 30|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
return openai_400() 31|if not f:
params_kw.llmid = recs[0].id 32| return openai_429()
params_kw.llmcatelogid = catelogid 33|
34|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 35|
if not f:
return openai_429()
return await inference(request, env=env)

View File

@ -1,34 +1,30 @@
# KTV Media API - subtitle-render 1|# KTV Media API - subtitle-render
# /v1/media/subtitle-render 2|# /v1/media/subtitle-render
3|
debug_params('params_kw', params_kw) 4|debug_params('params_kw', params_kw)
5|
userid = await get_user() 6|userid = await get_user()
userorgid = await get_userorgid() 7|userorgid = await get_userorgid()
if userid is None: 8|if userid is None:
return openai_403() 9| return openai_403()
10|
catelogid = 'ktv_pipeline' 11|catelogid = 'ktv_pipeline'
model_name = 'subtitle-render' 12|model_name = 'subtitle-render'
params_kw.model = model_name 13|params_kw.model = model_name
params_kw.catelogid = catelogid 14|params_kw.catelogid = catelogid
15|
env = request._run_ns 16|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 17|# llmid from cache (model+catelogid -> llmid)
sql = """select distinct a.* from llm a llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
join llm_api_map m on a.id = m.llmid if not llmid:
join llmcatelog b on m.llmcatelogid = b.id debug(f'model not found: params_kw.model')
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) return openai_400()
and a.model=${model}$ params_kw.llmid = llmid
and a.status = 'published'""" 28| params_kw.llmcatelogid = catelogid
recs = await sor.sqlExe(sql, {'catelogid': catelogid, 'model': model_name}) 29|
if len(recs) == 0: 30|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
return openai_400() 31|if not f:
params_kw.llmid = recs[0].id 32| return openai_429()
params_kw.llmcatelogid = catelogid 33|
34|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 35|
if not f:
return openai_429()
return await inference(request, env=env)

View File

@ -1,34 +1,30 @@
# KTV Media API - synth-generate 1|# KTV Media API - synth-generate
# /v1/media/synth-generate 2|# /v1/media/synth-generate
3|
debug_params('params_kw', params_kw) 4|debug_params('params_kw', params_kw)
5|
userid = await get_user() 6|userid = await get_user()
userorgid = await get_userorgid() 7|userorgid = await get_userorgid()
if userid is None: 8|if userid is None:
return openai_403() 9| return openai_403()
10|
catelogid = 'ktv_pipeline' 11|catelogid = 'ktv_pipeline'
model_name = 'synth-generate' 12|model_name = 'synth-generate'
params_kw.model = model_name 13|params_kw.model = model_name
params_kw.catelogid = catelogid 14|params_kw.catelogid = catelogid
15|
env = request._run_ns 16|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 17|# llmid from cache (model+catelogid -> llmid)
sql = """select distinct a.* from llm a llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
join llm_api_map m on a.id = m.llmid if not llmid:
join llmcatelog b on m.llmcatelogid = b.id debug(f'model not found: params_kw.model')
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) return openai_400()
and a.model=${model}$ params_kw.llmid = llmid
and a.status = 'published'""" 28| params_kw.llmcatelogid = catelogid
recs = await sor.sqlExe(sql, {'catelogid': catelogid, 'model': model_name}) 29|
if len(recs) == 0: 30|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
return openai_400() 31|if not f:
params_kw.llmid = recs[0].id 32| return openai_429()
params_kw.llmcatelogid = catelogid 33|
34|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 35|
if not f:
return openai_429()
return await inference(request, env=env)

View File

@ -1,34 +1,30 @@
# KTV Media API - video-eval-evaluate 1|# KTV Media API - video-eval-evaluate
# /v1/media/video-eval-evaluate 2|# /v1/media/video-eval-evaluate
3|
debug_params('params_kw', params_kw) 4|debug_params('params_kw', params_kw)
5|
userid = await get_user() 6|userid = await get_user()
userorgid = await get_userorgid() 7|userorgid = await get_userorgid()
if userid is None: 8|if userid is None:
return openai_403() 9| return openai_403()
10|
catelogid = 'ktv_pipeline' 11|catelogid = 'ktv_pipeline'
model_name = 'video-eval-evaluate' 12|model_name = 'video-eval-evaluate'
params_kw.model = model_name 13|params_kw.model = model_name
params_kw.catelogid = catelogid 14|params_kw.catelogid = catelogid
15|
env = request._run_ns 16|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 17|# llmid from cache (model+catelogid -> llmid)
sql = """select distinct a.* from llm a llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
join llm_api_map m on a.id = m.llmid if not llmid:
join llmcatelog b on m.llmcatelogid = b.id debug(f'model not found: params_kw.model')
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) return openai_400()
and a.model=${model}$ params_kw.llmid = llmid
and a.status = 'published'""" 28| params_kw.llmcatelogid = catelogid
recs = await sor.sqlExe(sql, {'catelogid': catelogid, 'model': model_name}) 29|
if len(recs) == 0: 30|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
return openai_400() 31|if not f:
params_kw.llmid = recs[0].id 32| return openai_429()
params_kw.llmcatelogid = catelogid 33|
34|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 35|
if not f:
return openai_429()
return await inference(request, env=env)

View File

@ -47,22 +47,11 @@ if not params_kw.lyrics:
lctype = params_kw.catelogid lctype = params_kw.catelogid
env = request._run_ns env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: llmid = await env.get_llmid_cached(env, params_kw.model, lctype)
# Look up llm by model name and catalog type through llm_api_map if not llmid:
sql = """select distinct a.* from llm a debug(f'{params_kw.model=} not found for catalog {lctype}')
join llm_api_map m on a.id = m.llmid return openai_400()
join llmcatelog b on m.llmcatelogid = b.id params_kw.llmid = llmid
where (b.id = ${lctype}$ OR b.name = ${lctype}$)
and a.name=${model}$
and a.status = 'published'"""
recs = await sor.sqlExe(sql, {
'lctype': lctype,
'model': params_kw.model
})
if len(recs) == 0:
debug(f'{params_kw.model=} not found for catalog {lctype}')
return openai_400()
params_kw.llmid = recs[0].id
params_kw.llmcatelogid = lctype params_kw.llmcatelogid = lctype
debug(f'{params_kw.llmid=}') debug(f'{params_kw.llmid=}')

View File

@ -1,74 +1,64 @@
# KTV Pipeline / GPU Service Inference API 1|# KTV Pipeline / GPU Service Inference API
# POST /v1/pipeline/submit 2|# POST /v1/pipeline/submit
# catelogid 固定为 ktv_pipeline无需传 catelogid 3|# catelogid 固定为 ktv_pipeline无需传 catelogid
# 4|#
# Required params: 5|# Required params:
# model: string - 模型名称,如 "ky-asr-transcribe" 6|# model: string - 模型名称,如 "ky-asr-transcribe"
# 7|#
# 各服务特定参数见 API 文档 8|# 各服务特定参数见 API 文档
# 9|#
# 同步服务直接返回结果: 10|# 同步服务直接返回结果:
# { 11|# {
# "taskid": "luid_xxx", 12|# "taskid": "luid_xxx",
# "taskstatus": "SUCCEEDED", 13|# "taskstatus": "SUCCEEDED",
# ...服务特定结果字段..., 14|# ...服务特定结果字段...,
# "usage": {...} 15|# "usage": {...}
# } 16|# }
# 17|#
# 异步服务返回任务信息: 18|# 异步服务返回任务信息:
# { 19|# {
# "taskid": "luid_xxx", 20|# "taskid": "luid_xxx",
# "taskstatus": "PENDING" 21|# "taskstatus": "PENDING"
# } 22|# }
# 异步结果通过 /v1/tasks?taskid=xxx 查询 23|# 异步结果通过 /v1/tasks?taskid=xxx 查询
24|
debug_params('params_kw', params_kw) 25|debug_params('params_kw', params_kw)
26|
userid = await get_user() 27|userid = await get_user()
userorgid = await get_userorgid() 28|userorgid = await get_userorgid()
if userid is None: 29|if userid is None:
debug('need login') 30| debug('need login')
return openai_403() 31| return openai_403()
32|
# Validate required parameters 33|# Validate required parameters
if not params_kw.model: 34|if not params_kw.model:
d = return_error('Missing required parameter: model') 35| d = return_error('Missing required parameter: model')
return json_response(d, status=400) 36| return json_response(d, status=400)
37|
catelogid = 'ktv_pipeline' 38|catelogid = 'ktv_pipeline'
params_kw.catelogid = catelogid 39|params_kw.catelogid = catelogid
40|
env = request._run_ns 41|env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: 42|# llmid from cache (model+catelogid -> llmid)
# Look up llm by model name through llm_api_map (fixed catelogid=ktv_pipeline) llmid = await env.get_llmid_cached(env, params_kw.model or 'qwen3-max', 'ktv_pipeline')
sql = """select distinct a.* from llm a if not llmid:
join llm_api_map m on a.id = m.llmid debug(f'model not found: params_kw.model')
join llmcatelog b on m.llmcatelogid = b.id return {{error_response}}
where (b.id = ${catelogid}$ OR b.name = ${catelogid}$) params_kw.llmid = llmid
and a.name=${model}$ 59| params_kw.llmcatelogid = catelogid
and a.status = 'published'""" 60|
recs = await sor.sqlExe(sql, { 61|debug(f'{params_kw.llmid=}')
'catelogid': catelogid, 62|
'model': params_kw.model 63|# Check balance
}) 64|f = await checkCustomerBalance(params_kw.llmid, userid, userorgid)
if len(recs) == 0: 65|if not f:
debug(f'{params_kw.model=} not found under ktv_pipeline catalog') 66| debug(f'{userid=} balance not enough')
d = return_error(f'Model "{params_kw.model}" not found or not available') 67| return openai_429()
return json_response(d, status=400) 68|
params_kw.llmid = recs[0].id 69|# Generate task ID
params_kw.llmcatelogid = catelogid 70|if not params_kw.transno:
71| params_kw.transno = getID()
debug(f'{params_kw.llmid=}') 72|
73|# Call inference
# Check balance 74|return await inference(request, env=env)
f = await checkCustomerBalance(params_kw.llmid, userid, userorgid) 75|
if not f:
debug(f'{userid=} balance not enough')
return openai_429()
# Generate task ID
if not params_kw.transno:
params_kw.transno = getID()
# Call inference
return await inference(request, env=env)

View File

@ -47,22 +47,11 @@ if not params_kw.prompt:
lctype = params_kw.catelogid lctype = params_kw.catelogid
env = request._run_ns env = request._run_ns
async with get_sor_context(env, 'llmage') as sor: llmid = await env.get_llmid_cached(env, params_kw.model, lctype)
# Look up llm by model name and catalog type through llm_api_map if not llmid:
sql = """select distinct a.* from llm a debug(f'{params_kw.model=} not found for catalog {lctype}')
join llm_api_map m on a.id = m.llmid return openai_400()
join llmcatelog b on m.llmcatelogid = b.id params_kw.llmid = llmid
where (b.id = ${lctype}$ OR b.name = ${lctype}$)
and a.name=${model}$
and a.status = 'published'"""
recs = await sor.sqlExe(sql, {
'lctype': lctype,
'model': params_kw.model
})
if len(recs) == 0:
debug(f'{params_kw.model=} not found for catalog {lctype}')
return openai_400()
params_kw.llmid = recs[0].id
params_kw.llmcatelogid = lctype params_kw.llmcatelogid = lctype
debug(f'{params_kw.llmid=}') debug(f'{params_kw.llmid=}')