Commit c4ef5508 authored by Yaowentong's avatar Yaowentong

视频转文字

parent 1236b9ff
import json
import time
import text_extract_service as service
class FakeRedis:
def __init__(self):
self.data = {}
def set(self, key, value):
self.data[key] = value
return True
def get(self, key):
return self.data.get(key)
class BrokenRedis:
def set(self, key, value):
raise TimeoutError("redis timeout")
def get(self, key):
raise TimeoutError("redis timeout")
def fake_extract_events_from_url(url):
yield sse_event({"message": {"url": url}, "type": 5})
yield sse_event({"message": "第一段", "type": 1})
yield sse_event({"message": "第二段", "type": 1})
def sse_event(data):
return f"data: {json.dumps(data, ensure_ascii=False)}\n\n".encode("utf-8")
def wait_task_done(client, task_id):
for _ in range(20):
response = client.get(f"/api/text-extract/tasks/{task_id}")
task = response.get_json()["data"]
if task["status"] in ("success", "failed"):
return task
time.sleep(0.05)
raise AssertionError("task did not finish in time")
def test_submit_multiple_urls_and_query_result(monkeypatch):
monkeypatch.setattr(service, "extract_events_from_url", fake_extract_events_from_url)
monkeypatch.setattr(service, "redis_client", FakeRedis())
service.tasks.clear()
client = service.app.test_client()
response = client.post(
"/api/text-extract/tasks",
json={
"urls": [
"https://example.com/video-a",
"https://example.com/video-b",
]
},
)
assert response.status_code == 202
body = response.get_json()
assert body["code"] == 200
assert len(body["data"]) == 2
assert body["data"][0]["task_id"] != body["data"][1]["task_id"]
task = wait_task_done(client, body["data"][0]["task_id"])
assert task["status"] == "success"
assert task["status_message"] == "文案提取完成"
assert task["text"] == "第一段第二段"
assert task["video_info"] == {"url": "https://example.com/video-a"}
assert "events" not in task
def test_get_result_by_query_task_id(monkeypatch):
monkeypatch.setattr(service, "extract_events_from_url", fake_extract_events_from_url)
monkeypatch.setattr(service, "redis_client", FakeRedis())
service.tasks.clear()
client = service.app.test_client()
submit_response = client.post(
"/api/text-extract/tasks",
json={"url": "https://example.com/video"},
)
task_id = submit_response.get_json()["data"][0]["task_id"]
wait_task_done(client, task_id)
response = client.get("/api/text-extract/result", query_string={"task_id": task_id})
assert response.status_code == 200
body = response.get_json()
assert body["data"]["task_id"] == task_id
assert body["data"]["status_message"] == "文案提取完成"
assert body["data"]["text"] == "第一段第二段"
assert body["data"]["video_info"] == {"url": "https://example.com/video"}
assert "events" not in body["data"]
def test_get_result_after_memory_cache_is_cleared(monkeypatch):
monkeypatch.setattr(service, "extract_events_from_url", fake_extract_events_from_url)
monkeypatch.setattr(service, "redis_client", FakeRedis())
service.tasks.clear()
client = service.app.test_client()
submit_response = client.post(
"/api/text-extract/tasks",
json={"url": "https://example.com/video"},
)
task_id = submit_response.get_json()["data"][0]["task_id"]
wait_task_done(client, task_id)
service.tasks.clear()
response = client.get(f"/api/text-extract/tasks/{task_id}")
assert response.status_code == 200
body = response.get_json()
assert body["data"]["task_id"] == task_id
assert body["data"]["text"] == "第一段第二段"
def test_submit_still_returns_when_redis_fails(monkeypatch):
monkeypatch.setattr(service, "extract_events_from_url", fake_extract_events_from_url)
monkeypatch.setattr(service, "redis_client", BrokenRedis())
service.tasks.clear()
client = service.app.test_client()
response = client.post(
"/api/text-extract/tasks",
json={"url": "https://example.com/video"},
)
assert response.status_code == 202
body = response.get_json()
assert body["code"] == 200
assert len(body["data"]) == 1
assert body["data"][0]["task_id"]
import json
import os
import sys
import threading
import traceback
import uuid
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime
from typing import Any, Dict, List, Optional
import redis
from flask import Flask, jsonify, request
base_path = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
sys.path.append(base_path)
app = Flask(__name__)
MAX_WORKERS = 4
REDIS_TASK_KEY_PREFIX = "text_extract:task:"
executor = ThreadPoolExecutor(max_workers=MAX_WORKERS)
redis_client = redis.Redis(
host="172.16.0.24",
port=6379,
password="aiyingli@@123",
socket_connect_timeout=2,
socket_timeout=10,
decode_responses=True,
db=0,
)
tasks: Dict[str, Dict[str, Any]] = {}
tasks_lock = threading.Lock()
def extract_events_from_url(url: str):
from ai_chat_stream import parse_utils
yield from parse_utils.extract_mp4_to_text(url)
def now_text() -> str:
return datetime.now().strftime("%Y-%m-%d %H:%M:%S")
def build_task(url: str) -> Dict[str, Any]:
return {
"task_id": str(uuid.uuid4()),
"url": url,
"status": "queued",
"status_message": "任务已提交",
"text": "",
"video_info": None,
"events": [],
"error": None,
"created_at": now_text(),
"updated_at": now_text(),
"finished_at": None,
}
def get_task_key(task_id: str) -> str:
return f"{REDIS_TASK_KEY_PREFIX}{task_id}"
def save_task_to_redis(task: Dict[str, Any]):
try:
redis_client.set(
get_task_key(task["task_id"]),
json.dumps(task, ensure_ascii=False),
)
except Exception as exc:
print(f"save task to redis failed, task_id={task.get('task_id')}, error={exc}")
def load_task_from_redis(task_id: str) -> Optional[Dict[str, Any]]:
try:
task_text = redis_client.get(get_task_key(task_id))
except Exception as exc:
print(f"load task from redis failed, task_id={task_id}, error={exc}")
return None
if not task_text:
return None
return json.loads(task_text)
def update_task(task_id: str, **kwargs):
task = None
with tasks_lock:
task = tasks.get(task_id)
if task is None:
task = load_task_from_redis(task_id)
if task is None:
return
with tasks_lock:
task.update(kwargs)
task["updated_at"] = now_text()
tasks[task_id] = task
save_task_to_redis(task)
def get_task(task_id: str) -> Optional[Dict[str, Any]]:
with tasks_lock:
task = tasks.get(task_id)
if task is None:
task = load_task_from_redis(task_id)
if task is None:
return None
with tasks_lock:
tasks[task_id] = task
return dict(task)
def build_public_task(task: Dict[str, Any]) -> Dict[str, Any]:
return {
"task_id": task.get("task_id"),
"url": task.get("url"),
"status": task.get("status"),
"status_message": task.get("status_message"),
"video_info": task.get("video_info"),
"text": task.get("text", ""),
"error": task.get("error"),
"created_at": task.get("created_at"),
"updated_at": task.get("updated_at"),
"finished_at": task.get("finished_at"),
}
def parse_sse_event(raw_event: Any) -> Optional[Dict[str, Any]]:
if isinstance(raw_event, bytes):
event_text = raw_event.decode("utf-8", errors="ignore")
else:
event_text = str(raw_event)
for line in event_text.splitlines():
line = line.strip()
if not line.startswith("data:"):
continue
data_text = line[len("data:"):].strip()
if not data_text or data_text == "[DONE]":
continue
try:
return json.loads(data_text)
except json.JSONDecodeError:
return {
"message": data_text,
"type": 0,
}
return None
def extract_text_from_url(task_id: str, url: str):
update_task(task_id, status="running", status_message="任务处理中")
events: List[Dict[str, Any]] = []
text_parts: List[str] = []
try:
for raw_event in extract_events_from_url(url):
event = parse_sse_event(raw_event)
if event is None:
continue
events.append(event)
if event.get("type") == 5 and isinstance(event.get("message"), dict):
update_task(
task_id,
video_info=event.get("message"),
status_message="已获取视频基础信息",
)
if event.get("type") == 3:
message = event.get("message")
if message is not None:
update_task(task_id, status_message=str(message))
if event.get("type") == 1:
message = event.get("message")
if message is not None:
text_parts.append(str(message))
update_task(
task_id,
events=list(events),
text="".join(text_parts),
)
final_text = "".join(text_parts)
if final_text:
update_task(
task_id,
status="success",
status_message="文案提取完成",
text=final_text,
events=events,
finished_at=now_text(),
)
else:
update_task(
task_id,
status="failed",
status_message="文案提取失败",
text="",
events=events,
error="未提取到文案",
finished_at=now_text(),
)
except Exception as exc:
update_task(
task_id,
status="failed",
status_message="文案提取异常",
error=f"{exc}\n{traceback.format_exc()}",
events=events,
text="".join(text_parts),
finished_at=now_text(),
)
def normalize_urls(payload: Dict[str, Any]) -> List[str]:
urls = payload.get("urls")
if urls is None:
url = payload.get("url")
urls = [url] if url else []
if not isinstance(urls, list):
raise ValueError("urls 必须是数组")
cleaned_urls = []
for item in urls:
if not isinstance(item, str):
raise ValueError("url 必须是字符串")
url = item.strip()
if url:
cleaned_urls.append(url)
if not cleaned_urls:
raise ValueError("请提交 url 或 urls")
return cleaned_urls
@app.route("/api/text-extract/tasks", methods=["POST"])
def submit_text_extract_tasks():
payload = request.get_json(silent=True) or {}
try:
urls = normalize_urls(payload)
except ValueError as exc:
return jsonify({
"code": 400,
"msg": str(exc),
"data": None,
}), 400
created_tasks = []
for url in urls:
task = build_task(url)
task_id = task["task_id"]
with tasks_lock:
tasks[task_id] = task
save_task_to_redis(task)
executor.submit(extract_text_from_url, task_id, url)
created_tasks.append({
"url": url,
"task_id": task_id,
"status": task["status"],
})
return jsonify({
"code": 200,
"msg": "success",
"data": created_tasks,
}), 202
@app.route("/api/text-extract/tasks/<task_id>", methods=["GET"])
def get_text_extract_task(task_id: str):
task = get_task(task_id)
if task is None:
return jsonify({
"code": 404,
"msg": "任务不存在",
"data": None,
}), 404
return jsonify({
"code": 200,
"msg": "success",
"data": build_public_task(task),
})
@app.route("/api/text-extract/result", methods=["GET"])
def get_text_extract_result_by_query():
task_id = request.args.get("task_id", "").strip()
if not task_id:
return jsonify({
"code": 400,
"msg": "请提供 task_id",
"data": None,
}), 400
return get_text_extract_task(task_id)
if __name__ == "__main__":
app.run(host="0.0.0.0", port=10443)
# 文案提取服务接口文档
## 服务说明
该服务用于提交视频/图文链接并异步提取文案。
服务启动后默认监听:
```text
http://172.16.18.10:10443
```
启动命令:
```bash
python3 text_extract_service.py
```
任务数据会写入 Redis:
```text
text_extract:task:{task_id}
```
因此服务重启后,只要 Redis 中仍有任务数据,就可以继续通过 `task_id` 查询结果。
## 接口列表
| 方法 | 路径 | 说明 |
| --- | --- | --- |
| POST | `/api/text-extract/tasks` | 提交一个或多个 URL,返回每个 URL 对应的任务 ID |
| GET | `/api/text-extract/tasks/{task_id}` | 根据任务 ID 查询文案提取结果 |
另提供一个等价查询接口:
```text
GET /api/text-extract/result?task_id={task_id}
```
## 1. 提交文案提取任务
### 请求
```http
POST /api/text-extract/tasks
Content-Type: application/json
```
### 请求参数
支持提交单个 URL:
```json
{
"url": "https://www.douyin.com/jingxuan?modal_id=7654800977366715698"
}
```
也支持一次提交多个 URL:
```json
{
"urls": [
"https://www.douyin.com/jingxuan?modal_id=7654800977366715698",
"https://www.xiaohongshu.com/discovery/item/xxxx?xsec_token=xxxx"
]
}
```
### 参数说明
| 字段 | 类型 | 必填 | 说明 |
| --- | --- | --- | --- |
| url | string | 否 | 单个待提取链接 |
| urls | array[string] | 否 | 多个待提取链接 |
说明:
- `url``urls` 二选一。
- 如果同时传入 `url``urls`,服务优先使用 `urls`
- 当前支持的平台取决于原文案提取逻辑,主要包括抖音、小红书、快手、B站。
### curl 示例
提交单个 URL:
```bash
curl -X POST "http://172.16.18.10:10443/api/text-extract/tasks" \
-H "Content-Type: application/json" \
-d '{
"url": "https://www.douyin.com/jingxuan?modal_id=7654800977366715698"
}'
```
提交多个 URL:
```bash
curl -X POST "http://172.16.18.10:10443/api/text-extract/tasks" \
-H "Content-Type: application/json" \
-d '{
"urls": [
"https://www.douyin.com/jingxuan?modal_id=7654800977366715698",
"https://www.xiaohongshu.com/discovery/item/xxxx?xsec_token=xxxx"
]
}'
```
### 成功响应
HTTP 状态码:
```text
202 Accepted
```
响应示例:
```json
{
"code": 200,
"msg": "success",
"data": [
{
"url": "https://www.douyin.com/jingxuan?modal_id=7654800977366715698",
"task_id": "0cb446fb-4938-4d52-a500-9d0e2fab109e",
"status": "queued"
}
]
}
```
### 响应字段说明
| 字段 | 类型 | 说明 |
| --- | --- | --- |
| code | number | 业务状态码,成功为 `200` |
| msg | string | 响应信息 |
| data | array | 任务列表 |
| data[].url | string | 提交的原始 URL |
| data[].task_id | string | 任务 ID,后续查询结果使用 |
| data[].status | string | 任务初始状态,通常为 `queued` |
### 失败响应
请求参数错误:
```json
{
"code": 400,
"msg": "请提交 url 或 urls",
"data": null
}
```
## 2. 查询文案提取结果
### 请求
```http
GET /api/text-extract/tasks/{task_id}
```
### 路径参数
| 字段 | 类型 | 必填 | 说明 |
| --- | --- | --- | --- |
| task_id | string | 是 | 提交任务时返回的任务 ID |
### curl 示例
```bash
curl -X GET "http://172.16.18.10:10443/api/text-extract/tasks/0cb446fb-4938-4d52-a500-9d0e2fab109e"
```
也可以使用 query 参数查询:
```bash
curl -X GET "http://172.16.18.10:10443/api/text-extract/result?task_id=0cb446fb-4938-4d52-a500-9d0e2fab109e"
```
### 成功响应
任务处理中:
```json
{
"code": 200,
"msg": "success",
"data": {
"task_id": "0cb446fb-4938-4d52-a500-9d0e2fab109e",
"url": "https://www.douyin.com/jingxuan?modal_id=7654800977366715698",
"status": "running",
"status_message": "视频文案提取中",
"video_info": {
"aweme_type": "SHIPIN",
"title": "视频标题",
"item_title": "视频标题",
"cover": "https://static2.douchacha.com/assets/xxx/cover.jpg",
"duration": "268284",
"comment_count": "386",
"digg_count": "6123",
"share_count": "870",
"video_id": "7654800977366715698",
"platform": "DY"
},
"text": "",
"error": null,
"created_at": "2026-07-07 20:02:12",
"updated_at": "2026-07-07 20:02:30",
"finished_at": null
}
}
```
任务成功:
```json
{
"code": 200,
"msg": "success",
"data": {
"task_id": "0cb446fb-4938-4d52-a500-9d0e2fab109e",
"url": "https://www.douyin.com/jingxuan?modal_id=7654800977366715698",
"status": "success",
"status_message": "文案提取完成",
"video_info": {
"aweme_type": "SHIPIN",
"title": "视频标题",
"item_title": "视频标题",
"cover": "https://static2.douchacha.com/assets/xxx/cover.jpg",
"duration": "268284",
"comment_count": "386",
"digg_count": "6123",
"share_count": "870",
"video_id": "7654800977366715698",
"platform": "DY"
},
"text": "这里是最终提取出来的视频文案",
"error": null,
"created_at": "2026-07-07 20:02:12",
"updated_at": "2026-07-07 20:02:59",
"finished_at": "2026-07-07 20:02:59"
}
}
```
任务失败:
```json
{
"code": 200,
"msg": "success",
"data": {
"task_id": "0cb446fb-4938-4d52-a500-9d0e2fab109e",
"url": "https://example.com/wrong-url",
"status": "failed",
"status_message": "文案提取失败",
"video_info": null,
"text": "",
"error": "未提取到文案",
"created_at": "2026-07-07 20:02:12",
"updated_at": "2026-07-07 20:02:59",
"finished_at": "2026-07-07 20:02:59"
}
}
```
### 响应字段说明
| 字段 | 类型 | 说明 |
| --- | --- | --- |
| code | number | 业务状态码,成功为 `200` |
| msg | string | 响应信息 |
| data.task_id | string | 任务 ID |
| data.url | string | 提交的原始 URL |
| data.status | string | 任务状态 |
| data.status_message | string | 当前任务阶段说明 |
| data.video_info | object/null | 视频基础信息 |
| data.text | string | 提取出的文案结果 |
| data.error | string/null | 失败或异常信息 |
| data.created_at | string | 任务创建时间 |
| data.updated_at | string | 最近更新时间 |
| data.finished_at | string/null | 任务完成时间 |
## 任务状态说明
| 状态 | 说明 |
| --- | --- |
| queued | 任务已提交,等待后台线程处理 |
| running | 任务处理中 |
| success | 文案提取成功 |
| failed | 文案提取失败 |
## status_message 常见值
| status_message | 说明 |
| --- | --- |
| 任务已提交 | 已创建任务 |
| 任务处理中 | 后台线程已开始处理 |
| 已获取视频基础信息 | 已解析到视频基础信息 |
| 视频文案提取中 | 正在提取视频文案 |
| 识别完成 | ASR 识别完成 |
| 文案提取完成 | 已得到最终文案 |
| 文案提取失败 | 任务执行完成,但未提取到文案 |
| 文案提取异常 | 任务执行过程中出现异常 |
## 查询建议
提交任务后可以每 2-3 秒查询一次结果。
示例:
```bash
curl -X GET "http://172.16.18.10:10443/api/text-extract/tasks/你的task_id"
```
当返回:
```json
{
"status": "success"
}
```
即可读取 `data.text` 作为最终文案结果。
## 注意事项
- `POST` 接口只负责提交任务,不会等待文案提取完成。
- 文案提取是后台异步执行的,耗时取决于爬虫解析、TOS 文件生成、ASR 识别等外部服务。
- 查询接口不会返回内部 `events`,只返回视频基础信息 `video_info` 和文案结果 `text`
- Redis 不可用时,当前进程内仍可查询内存中的任务;但服务重启后无法恢复未写入 Redis 的任务。
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment