For agentic workers: REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (
- [ ]) syntax for tracking.
Goal: 在 cfc-langgraph/ 建立 Python LangGraph sidecar 服务,上线 RecommendAgent,Java 端支持灰度切换(5% 流量走 Python,异常自动 fallback)。
Architecture: FastAPI 服务 (port 9000) 承载 LangChain Agent,通过 HTTP 与 Java (port 9082) 通信。数据存储用 ChromaDB 文件模式 + Java 侧 MySQL。Agent 调 LLM API(DeepSeek-V3)进行推荐+解释。
Tech Stack: Python 3.11, FastAPI, LangChain, LangGraph, ChromaDB, Docker, Java 8 (不变)
D:\workspace\cfc\Files:
cfc-langgraph/pyproject.tomlcfc-langgraph/.env.examplecfc-langgraph/.gitignorecfc-langgraph/Dockerfilecfc-langgraph/app/__init__.pycfc-langgraph/app/main.pyInterfaces:
Produces: FastAPI app 入口 app.main:app,/health 返回 200
[x] Step 1: 创建 pyproject.toml
[project]
name = "cfc-langgraph"
version = "0.1.0"
description = "cfc AI sidecar service with LangGraph"
requires-python = ">=3.11"
dependencies = [
"fastapi>=0.115",
"uvicorn[standard]>=0.34",
"pydantic>=2.10",
"pydantic-settings>=2.7",
"httpx>=0.28",
"langchain>=0.3",
"langchain-community>=0.3",
"langchain-openai>=0.3",
"langgraph>=0.3",
"chromadb>=0.6",
"langchain-chroma>=0.2",
"python-multipart>=0.0.20",
]
[project.optional-dependencies]
dev = [
"pytest>=8",
"pytest-asyncio>=0.25",
"ruff>=0.8",
"httpx",
]
[tool.ruff]
target-version = "py311"
line-length = 100
select = ["E", "F", "I", "N", "W"]
ignore = ["E501"]
[x] Step 2: 创建 .env.example
# LLM
LLM_API_KEY=sk-your-key
LLM_BASE_URL=https://api.deepseek.com/v1
LLM_MODEL=deepseek-chat
# Embedding
EMBEDDING_API_KEY=${LLM_API_KEY}
EMBEDDING_BASE_URL=https://api.deepseek.com/v1
EMBEDDING_MODEL=text-embedding-v3
# Java Backend
JAVA_BASE_URL=http://localhost:9082
JAVA_CONTEXT_URL=${JAVA_BASE_URL}/api/ai/context
# LangSmith (optional)
LANGCHAIN_TRACING_V2=false
LANGCHAIN_API_KEY=
LANGCHAIN_PROJECT=cfc-langgraph
# Service
SERVICE_HOST=0.0.0.0
SERVICE_PORT=9000
LOG_LEVEL=info
# Chroma
CHROMA_DB_PATH=./data/chroma_db
[x] Step 3: 创建 .gitignore
__pycache__/
*.py[cod]
.env
data/
*.egg-info/
dist/
.ruff_cache/
.pytest_cache/
[x] Step 4: 创建 Dockerfile
FROM python:3.11-slim
WORKDIR /app
COPY pyproject.toml .
RUN pip install --no-cache-dir -e .
COPY app/ app/
COPY data/ data/ 2>/dev/null || true
EXPOSE 9000
CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "9000"]
[x] Step 5: 创建 app/__init__.py(空文件)
[x] Step 6: 创建 app/main.py
from fastapi import FastAPI
from app.api import health
app = FastAPI(title="cfc-langgraph", version="0.1.0")
app.include_router(health.router)
@app.on_event("startup")
async def startup():
pass # 后续 Phase 在此初始化 RAG / Agent
@app.on_event("shutdown")
async def shutdown():
pass # 后续 Phase 在此清理资源
[x] Step 7: 验证服务可启动
cd cfc-langgraph
pip install -e .
uvicorn app.main:app --port 9000 &
curl http://localhost:9000/health
# 预期: {"status":"ok"}
kill %1
[x] Step 8: 创建 app/api/__init__.py(空文件)
[x] Step 9: 创建 app/api/health.py
from fastapi import APIRouter
router = APIRouter(tags=["health"])
@router.get("/health")
async def health():
return {"status": "ok"}
[x] Step 10: Commit
git add cfc-langgraph/
git commit -m "feat(langgraph): scaffold Python sidecar project"
Files:
cfc-langgraph/app/config.pyInterfaces:
Produces: app.config.settings — Pydantic Settings 单例,所有模块通过 from app.config import settings 读取
[x] Step 1: 创建 app/config.py
from pydantic_settings import BaseSettings
from typing import Optional
class Settings(BaseSettings):
# LLM
llm_api_key: str
llm_base_url: str = "https://api.deepseek.com/v1"
llm_model: str = "deepseek-chat"
# Embedding
embedding_api_key: Optional[str] = None
embedding_base_url: Optional[str] = None
embedding_model: str = "text-embedding-v3"
# Java Backend
java_base_url: str = "http://localhost:9082"
java_context_url: Optional[str] = None
# LangSmith
langchain_tracing_v2: bool = False
langchain_api_key: Optional[str] = None
langchain_project: str = "cfc-langgraph"
# Service
service_host: str = "0.0.0.0"
service_port: int = 9000
log_level: str = "info"
# Chroma
chroma_db_path: str = "./data/chroma_db"
model_config = {"env_file": ".env", "env_file_encoding": "utf-8"}
@property
def effective_embedding_api_key(self) -> str:
return self.embedding_api_key or self.llm_api_key
@property
def effective_embedding_base_url(self) -> str:
return self.embedding_base_url or self.llm_base_url
@property
def effective_java_context_url(self) -> str:
return self.java_context_url or f"{self.java_base_url}/api/ai/context"
settings = Settings()
[x] Step 2: 创建 .env 文件(从 .env.example 复制,填入真实 Key)
cp cfc-langgraph/.env.example cfc-langgraph/.env
# 手动编辑填入 LLM_API_KEY
[x] Step 3: 验证配置可加载
cd cfc-langgraph
python -c "from app.config import settings; print(settings.llm_model)"
# 预期输出: deepseek-chat
[x] Step 4: Commit
git add cfc-langgraph/
git commit -m "feat(langgraph): config management with pydantic-settings"
Files:
cfc-langgraph/app/models/__init__.pycfc-langgraph/app/models/common.pycfc-langgraph/app/models/chat.pycfc-langgraph/app/models/recommend.pyInterfaces:
Produces: Phase 1 所需的请求/响应模型
[x] Step 1: 创建 app/models/__init__.py
from .common import UserContext, SourceInfo
from .chat import ChatRequest, ChatResponse
from .recommend import RecommendRequest, RecommendResponse, RecommendItem
[x] Step 2: 创建 app/models/common.py
from pydantic import BaseModel
from typing import Optional
class UserContext(BaseModel):
"""Java 侧传来的业务上下文"""
child_id: Optional[int] = None
report_id: Optional[int] = None
survey_id: Optional[int] = None
family_id: Optional[int] = None
mascot_code: Optional[str] = None
class SourceInfo(BaseModel):
"""回答引用来源"""
type: str # knowledge / tool
title: str
score: Optional[float] = None
[x] Step 3: 创建 app/models/chat.py
from pydantic import BaseModel
from typing import Optional
from .common import UserContext, SourceInfo
class ChatRequest(BaseModel):
query: str
user_id: int
conversation_id: str = ""
context: Optional[UserContext] = None
class ChatResponse(BaseModel):
answer: str
conversation_id: str
sources: list[SourceInfo] = []
tasks: list[dict] = []
trace_id: str = ""
[x] Step 4: 创建 app/models/recommend.py
from pydantic import BaseModel
from typing import Optional
class RecommendItem(BaseModel):
type: str # product / activity / article
id: int
name: str
description: str = ""
cover_image: str = ""
price: Optional[float] = None
url: str = ""
score: float = 0.0
reason: str = ""
class RecommendRequest(BaseModel):
user_id: int
query: str = ""
tags: list[str] = []
types: Optional[list[str]] = None
limit: int = 5
context: Optional[dict] = None
class RecommendResponse(BaseModel):
items: list[RecommendItem]
source: str = "" # sql / vector / hybrid
trace_id: str = ""
[x] Step 5: Commit
git add cfc-langgraph/
git commit -m "feat(langgraph): pydantic models for API"
Files:
cfc-langgraph/app/rag/__init__.pycfc-langgraph/app/rag/embeddings.pycfc-langgraph/app/rag/retriever.pyInterfaces:
embeddings.get_embeddings() → Embeddings 实例Produces: retriever.RagRetriever 类,提供 retrieve(query, filters, k) 方法
[x] Step 1: 创建 app/rag/__init__.py
from .embeddings import get_embeddings
from .retriever import RagRetriever
[x] Step 2: 创建 app/rag/embeddings.py
from langchain_openai import OpenAIEmbeddings
from app.config import settings
_embeddings = None
def get_embeddings():
global _embeddings
if _embeddings is None:
_embeddings = OpenAIEmbeddings(
model=settings.embedding_model,
api_key=settings.effective_embedding_api_key,
base_url=settings.effective_embedding_base_url,
)
return _embeddings
[x] Step 3: 创建 app/rag/retriever.py
from langchain_chroma import Chroma
from langchain.retrievers import EnsembleRetriever
from langchain_community.retrievers import BM25Retriever
from langchain.retrievers.document_compressors import LLMChainExtractor
from langchain.retrievers import ContextualCompressionRetriever
from .embeddings import get_embeddings
from app.config import settings
from app.tools.java_client import JavaClient
from typing import Optional
import logging
logger = logging.getLogger(__name__)
class RagRetriever:
"""混合检索器: ChromaDB 向量 + BM25 关键词 + 可选 LLM 压缩"""
def __init__(self, collection_name: str = "cfc_knowledge"):
embeddings = get_embeddings()
self.vectorstore = Chroma(
collection_name=collection_name,
embedding_function=embeddings,
persist_directory=settings.chroma_db_path,
)
self.java_client = JavaClient()
self._bm25_retriever = None
self._bm25_texts = []
async def initialize(self):
"""从 Java 侧拉取知识库数据, 构建 BM25 索引"""
try:
articles = await self.java_client.get_published_articles()
self._bm25_texts = [
f"{a['title']} {a['summary']} {a.get('tags', '')}"
for a in articles
]
if self._bm25_texts:
self._bm25_retriever = BM25Retriever.from_texts(
self._bm25_texts, metadatas=articles
)
logger.info("BM25 索引就绪: %d 条", len(self._bm25_texts))
except Exception as e:
logger.warning("BM25 初始化失败(不影响向量检索): %s", e)
async def retrieve(
self,
query: str,
filters: Optional[dict] = None,
k: int = 5,
use_compression: bool = False,
) -> list[dict]:
"""混合检索, 返回 [{content, metadata, score}]"""
retrievers = []
# 向量检索
vector_retriever = self.vectorstore.as_retriever(
search_kwargs={"k": k, "filter": filters}
)
retrievers.append(vector_retriever)
# BM25 检索
if self._bm25_retriever:
retrievers.append(self._bm25_retriever)
if len(retrievers) == 1:
docs = await retrievers[0].ainvoke(query)
else:
ensemble = EnsembleRetriever(
retrievers=retrievers, weights=[0.6, 0.4]
)
docs = await ensemble.ainvoke(query)
# 可选: LLM 压缩去噪
if use_compression and docs:
from langchain_openai import ChatOpenAI
llm = ChatOpenAI(
model=settings.llm_model,
api_key=settings.llm_api_key,
base_url=settings.llm_base_url,
)
compressor = LLMChainExtractor.from_llm(llm)
compression_retriever = ContextualCompressionRetriever(
base_compressor=compressor,
base_retriever=self.vectorstore.as_retriever(),
)
docs = await compression_retriever.ainvoke(query)
results = []
for doc in docs:
results.append({
"content": doc.page_content,
"metadata": doc.metadata,
"score": getattr(doc, "metadata", {}).get("score", 0),
})
return results[:k]
[x] Step 4: Commit
git add cfc-langgraph/
git commit -m "feat(langgraph): ChromaDB + RAG retriever infrastructure"
Files:
cfc-langgraph/app/tools/__init__.pycfc-langgraph/app/tools/java_client.pyInterfaces:
JavaClient 类, 封装对所有 Java 接口的 HTTP 调用Produces: get_published_articles(), search_products(), search_activities(), search_articles()
[x] Step 1: 创建 app/tools/__init__.py
from .java_client import JavaClient
[x] Step 2: 创建 app/tools/java_client.py
import httpx
from typing import Optional
from app.config import settings
import logging
logger = logging.getLogger(__name__)
class JavaClient:
"""Java 后端 HTTP 客户端 (所有 Python→Java 通信的单一入口)"""
def __init__(self):
self.base_url = settings.java_base_url
self._client: Optional[httpx.AsyncClient] = None
async def _get_client(self) -> httpx.AsyncClient:
if self._client is None:
self._client = httpx.AsyncClient(
base_url=self.base_url,
timeout=httpx.Timeout(10.0, connect=3.0),
)
return self._client
async def close(self):
if self._client:
await self._client.aclose()
self._client = None
async def get_published_articles(self) -> list[dict]:
"""获取已发布的文章列表 (用于构建知识库)"""
client = await self._get_client()
resp = await client.post("/api/article/list", json={"status": "published", "limit": 1000})
data = resp.json()
if data.get("code") == 200:
return data.get("data", [])
return []
async def search_products(self, keyword: str, limit: int = 5) -> list[dict]:
"""按关键词搜索上架商品"""
client = await self._get_client()
resp = await client.post("/api/product/search", json={
"keyword": keyword, "status": "上架", "limit": limit
})
data = resp.json()
if data.get("code") == 200:
return data.get("data", [])
return []
async def search_activities(self, keyword: str, limit: int = 5) -> list[dict]:
"""按关键词搜索进行中的活动"""
client = await self._get_client()
resp = await client.post("/api/activity/search", json={
"keyword": keyword, "status": "published", "limit": limit
})
data = resp.json()
if data.get("code") == 200:
return data.get("data", [])
return []
async def search_articles(self, keyword: str, limit: int = 5) -> list[dict]:
"""按关键词搜索已发布文章"""
client = await self._get_client()
resp = await client.post("/api/article/search", json={
"keyword": keyword, "status": "published", "limit": limit
})
data = resp.json()
if data.get("code") == 200:
return data.get("data", [])
return []
async def get_user_context(self, user_id: int, params: Optional[dict] = None) -> dict:
"""获取用户上下文 (对应 Java AiContextService)"""
client = await self._get_client()
resp = await client.post(settings.effective_java_context_url, json={
"user_id": str(user_id),
"params": params or {},
})
data = resp.json()
if data.get("code") == 200:
return data.get("data", {})
return {}
[x] Step 3: Commit
git add cfc-langgraph/
git commit -m "feat(langgraph): Java HTTP client for tool calls"
Files:
cfc-langgraph/app/tools/product_tools.pycfc-langgraph/app/agents/__init__.pycfc-langgraph/app/agents/recommend_agent.pycfc-langgraph/app/graphs/__init__.pycfc-langgraph/app/graphs/recommend_graph.pycfc-langgraph/app/api/recommend.pycfc-langgraph/app/main.py (注册路由)Interfaces:
JavaClient, RagRetriever, ChatOpenAIProduces: POST /api/v1/recommend 接口
[x] Step 1: 创建 app/tools/product_tools.py
from langchain_core.tools import tool
from app.tools.java_client import JavaClient
import logging
logger = logging.getLogger(__name__)
_java = JavaClient()
@tool
async def search_product_by_keyword(keyword: str, limit: int = 5) -> str:
"""按关键词搜索上架商品, 返回 JSON 商品列表"""
try:
products = await _java.search_products(keyword, limit)
if not products:
return "[]"
return str([{
"id": p["id"],
"name": p["name"],
"price": p.get("price"),
"description": p.get("intro") or p.get("description", ""),
} for p in products])
except Exception as e:
logger.warning("搜索商品失败: %s", e)
return "[]"
@tool
async def search_activity_by_keyword(keyword: str, limit: int = 5) -> str:
"""按关键词搜索进行中的活动, 返回 JSON 活动列表"""
try:
activities = await _java.search_activities(keyword, limit)
if not activities:
return "[]"
return str([{
"id": a["id"],
"name": a["title"],
"description": a.get("description", ""),
} for a in activities])
except Exception as e:
logger.warning("搜索活动失败: %s", e)
return "[]"
@tool
async def search_article_by_keyword(keyword: str, limit: int = 5) -> str:
"""按关键词搜索已发布文章, 返回 JSON 文章列表"""
try:
articles = await _java.search_articles(keyword, limit)
if not articles:
return "[]"
return str([{
"id": a["id"],
"name": a["title"],
"summary": a.get("summary", ""),
} for a in articles])
except Exception as e:
logger.warning("搜索文章失败: %s", e)
return "[]"
[x] Step 2: 创建 app/agents/__init__.py(空)
[x] Step 3: 创建 app/graphs/__init__.py(空)
[x] Step 4: 创建 app/graphs/recommend_graph.py
Phase 1 用最简单的单步 Agent(不涉及 StateGraph 复杂特性),后续 Phase 再升级为图编排:
from langchain_openai import ChatOpenAI
from langchain_core.messages import SystemMessage, HumanMessage
from app.tools.product_tools import search_product_by_keyword, search_activity_by_keyword, search_article_by_keyword
from app.config import settings
import json
import logging
logger = logging.getLogger(__name__)
SYSTEM_PROMPT = """你是一个儿童成长营养推荐助手。根据用户的需求和营养标签, 推荐合适的商品、活动或文章。
推荐原则:
1. 首先尝试使用搜索工具查找匹配的内容
2. 如果搜索结果为空, 基于你的知识给出建议
3. 每项推荐必须附带推荐理由
4. 以 JSON 格式输出推荐结果
输出格式:
{
"items": [
{
"source": "tool" 或 "knowledge",
"type": "product" / "activity" / "article",
"id": 数字,
"name": "名称",
"description": "描述",
"reason": "为什么推荐这个"
}
]
}
"""
class RecommendAgent:
def __init__(self):
self.llm = ChatOpenAI(
model=settings.llm_model,
api_key=settings.llm_api_key,
base_url=settings.llm_base_url,
temperature=0.3,
)
self.tools = [
search_product_by_keyword,
search_activity_by_keyword,
search_article_by_keyword,
]
self.llm_with_tools = self.llm.bind_tools(self.tools)
async def run(self, query: str, tags: list[str], limit: int = 5) -> dict:
"""执行推荐 Agent, 返回推荐结果"""
# 如果传入了 tags, 构造搜索关键词
search_query = query or " ".join(tags)
messages = [
SystemMessage(content=SYSTEM_PROMPT),
HumanMessage(content=f"用户需求: {search_query}\n最大返回数量: {limit}\n请搜索并推荐合适的内容。"),
]
# LangChain Tool calling 自动完成: LLM 决定调哪个 Tool → 工具返回结果 → LLM 组织回答
response = await self.llm_with_tools.ainvoke(messages)
# 尝试解析 JSON 输出
content = response.content
try:
# 提取 JSON 块
if "```json" in content:
json_str = content.split("```json")[1].split("```")[0].strip()
elif "```" in content:
json_str = content.split("```")[1].split("```")[0].strip()
else:
json_str = content.strip()
result = json.loads(json_str)
return result
except (json.JSONDecodeError, IndexError):
# 非 JSON 输出, 包装为文本回答
logger.warning("Agent 输出非 JSON, raw: %s", content[:200])
return {"items": [], "text": content}
[x] Step 5: 创建 app/api/recommend.py
from fastapi import APIRouter
from app.models.recommend import RecommendRequest, RecommendResponse, RecommendItem
from app.graphs.recommend_graph import RecommendAgent
from app.rag.retriever import RagRetriever
import logging
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/v1", tags=["recommend"])
_agent: RecommendAgent = None
_retriever: RagRetriever = None
def get_agent() -> RecommendAgent:
global _agent
if _agent is None:
_agent = RecommendAgent()
return _agent
def get_retriever() -> RagRetriever:
global _retriever
if _retriever is None:
_retriever = RagRetriever()
return _retriever
@router.post("/recommend", response_model=RecommendResponse)
async def recommend(req: RecommendRequest):
"""营养推荐: Agent 搜索+LLM 解释"""
try:
agent = get_agent()
result = await agent.run(query=req.query, tags=req.tags, limit=req.limit)
items_data = result.get("items", [])
items = []
for item in items_data:
items.append(RecommendItem(
type=item.get("type", "product"),
id=item.get("id", 0),
name=item.get("name", ""),
description=item.get("description", ""),
reason=item.get("reason", ""),
))
return RecommendResponse(items=items, source="agent")
except Exception as e:
logger.error("RecommendAgent 调用失败: %s", e, exc_info=True)
return RecommendResponse(items=[], source="error")
[x] Step 6: 修改 app/main.py 注册路由
from fastapi import FastAPI
from app.api import health, recommend
app = FastAPI(title="cfc-langgraph", version="0.1.0")
app.include_router(health.router)
app.include_router(recommend.router)
@app.on_event("startup")
async def startup():
from app.rag.retriever import RagRetriever
retriever = RagRetriever()
await retriever.initialize()
@app.on_event("shutdown")
async def shutdown():
from app.tools.java_client import JavaClient
client = JavaClient()
await client.close()
[x] Step 7: 验证 Recommend API
cd cfc-langgraph
# 确保 Java 后端在运行 (localhost:9082)
uvicorn app.main:app --port 9000 &
curl -X POST http://localhost:9000/api/v1/recommend \
-H "Content-Type: application/json" \
-d '{"user_id": 1, "tags": ["益生菌", "肠胃"], "limit": 3}'
# 预期: {"items": [...], "source": "agent"}
kill %1
[x] Step 8: Commit
git add cfc-langgraph/
git commit -m "feat(langgraph): RecommendAgent with tool calling"
Files:
cfc-backend/src/main/java/com/etotem/cfc/service/AiGateway.javacfc-backend/src/main/resources/application.yml(加 python.* 配置)Interfaces:
AiGateway.recommend() — Java 端统一的 Python 调用入口,超时/异常自动 fallbackConsumes: Python POST /api/v1/recommend
[x] Step 1: 在 application.yml 中新增配置
# application.yml 底部追加
python:
enabled: true
base-url: http://localhost:9000
timeout-ms: 15000
circuit-breaker:
failure-threshold: 3
reset-timeout-ms: 30000
[x] Step 2: 创建 AiGateway.java
package com.etotem.cfc.service;
import com.etotem.cfc.dto.RecommendationResult;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.http.HttpEntity;
import org.springframework.http.ResponseEntity;
import org.springframework.stereotype.Service;
import org.springframework.web.client.RestTemplate;
import javax.annotation.PostConstruct;
import java.util.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
/**
* LangGraph Python 服务网关
* 统一管理 Python 服务调用、熔断、Fallback
*/
@Service
public class AiGateway {
private static final Logger log = LoggerFactory.getLogger(AiGateway.class);
@Value("${python.enabled:false}")
private boolean enabled;
@Value("${python.base-url:http://localhost:9000}")
private String baseUrl;
@Value("${python.timeout-ms:15000}")
private int timeoutMs;
@Value("${python.circuit-breaker.failure-threshold:3}")
private int failureThreshold;
@Value("${python.circuit-breaker.reset-timeout-ms:30000}")
private int resetTimeoutMs;
private final RestTemplate restTemplate = new RestTemplate();
private final ObjectMapper objectMapper = new ObjectMapper();
// 熔断器状态
private final AtomicInteger consecutiveFailures = new AtomicInteger(0);
private final AtomicLong lastFailureTime = new AtomicLong(0);
private volatile boolean circuitOpen = false;
@PostConstruct
public void init() {
if (enabled) {
log.info("AiGateway 已启用: baseUrl={}, timeout={}ms", baseUrl, timeoutMs);
} else {
log.info("AiGateway 已禁用, 所有请求走 Dify");
}
}
public boolean isEnabled() {
return enabled;
}
private boolean isCircuitOpen() {
if (!circuitOpen) return false;
// 检查是否到重置时间
if (System.currentTimeMillis() - lastFailureTime.get() > resetTimeoutMs) {
circuitOpen = false;
consecutiveFailures.set(0);
log.info("AiGateway 熔断器已半开");
return false;
}
return true;
}
private void recordFailure() {
int failures = consecutiveFailures.incrementAndGet();
lastFailureTime.set(System.currentTimeMillis());
if (failures >= failureThreshold) {
circuitOpen = true;
log.warn("AiGateway 熔断器已打开 (连续{}次失败)", failures);
}
}
/**
* 调用 Python 推荐服务, 失败时返回 null (由调用方决定 fallback)
*/
public List<RecommendationResult> recommend(Long userId, List<String> tags, List<String> types, Integer limit) {
if (!enabled || isCircuitOpen()) return null;
try {
// 使用 ObjectMapper 构造 JSON 替代手动拼接
ObjectNode body = objectMapper.createObjectNode();
body.put("user_id", userId);
if (tags != null && !tags.isEmpty()) {
ArrayNode tagsArray = body.putArray("tags");
tags.forEach(tagsArray::add);
}
if (types != null && !types.isEmpty()) {
ArrayNode typesArray = body.putArray("types");
types.forEach(typesArray::add);
}
body.put("limit", limit != null ? limit : 5);
HttpEntity<String> entity = new HttpEntity<>(body.toString(), createJsonHeaders());
String url = baseUrl + "/api/v1/recommend";
ResponseEntity<String> response = restTemplate.postForEntity(url, entity, String.class);
if (response.getStatusCode().is2xxSuccessful() && response.getBody() != null) {
JsonNode root = objectMapper.readTree(response.getBody());
JsonNode items = root.get("items");
if (items != null && items.isArray()) {
List<RecommendationResult> results = new ArrayList<>();
for (JsonNode item : items) {
RecommendationResult r = new RecommendationResult();
r.setType(item.get("type").asText());
r.setId(item.get("id").asLong());
r.setName(item.get("name").asText());
r.setDescription(item.has("description") ? item.get("description").asText() : "");
r.setReason(item.has("reason") ? item.get("reason").asText() : "");
r.setSource("langgraph");
results.add(r);
}
consecutiveFailures.set(0);
log.debug("AiGateway recommend 成功: {} 条", results.size());
return results;
}
}
return null;
} catch (Exception e) {
log.warn("AiGateway recommend 调用失败: {}", e.getMessage());
recordFailure();
return null;
}
}
private org.springframework.http.HttpHeaders createJsonHeaders() {
org.springframework.http.HttpHeaders headers = new org.springframework.http.HttpHeaders();
headers.setContentType(org.springframework.http.MediaType.APPLICATION_JSON);
return headers;
}
}
[x] Step 3: 编译验证
cd cfc-backend
mvn clean compile -q
# 预期: BUILD SUCCESS
[x] Step 4: Commit
git add cfc-backend/src/main/java/com/etotem/cfc/service/AiGateway.java
git add cfc-backend/src/main/resources/application.yml
git commit -m "feat(backend): AiGateway with circuit breaker for LangGraph"
Files:
cfc-backend/src/main/java/com/etotem/cfc/controller/ai/RecommendationController.javacfc-backend/src/main/java/com/etotem/cfc/service/RecommendationService.javaInterfaces:
AiGateway.recommend()灰度策略: 5% 流量走 Python, 异常自动 fallback 现有逻辑
[x] Step 1: 修改 RecommendationController.java
// 注入 AiGateway
@Resource
private AiGateway aiGateway;
@Operation(summary = "按营养标签搜索推荐内容")
@PostMapping("/search")
public Result<List<RecommendationResult>> search(
@RequestBody RecommendationQuery query,
HttpServletRequest request) {
Long userId = (Long) request.getAttribute("userId");
if (userId == null) userId = 0L;
// 灰度: 5% 流量走 Python LangGraph
if (aiGateway.isEnabled() && Math.random() < 0.05) {
List<RecommendationResult> pythonResults = aiGateway.recommend(
userId, query.getNutritionTags(), query.getTypes(), query.getLimit());
if (pythonResults != null && !pythonResults.isEmpty()) {
log.info("[灰度] Recommend via LangGraph, userId={}", userId);
return Result.success(pythonResults);
}
log.info("[灰度] LangGraph 无结果/异常, fallback 本地搜索");
}
// 现有业务逻辑
List<RecommendationResult> results = recommendationService.search(query);
return Result.success(results);
}
注意: Math.random() < 0.05 需要在 Controller 中注入 HttpServletRequest,当前方法签名没有这个参数。Java 端统一通过 @RequestAttribute("userId") 获取,所以需要修改方法签名或改用 RequestContextHolder。
安全且兼容的改法:
@PostMapping("/search")
public Result<List<RecommendationResult>> search(
@RequestBody RecommendationQuery query,
@RequestAttribute(value = "userId", required = false) Long userId) {
// ...
}
[x] Step 2: 编译验证
cd cfc-backend
mvn clean compile -q
# 预期: BUILD SUCCESS
[x] Step 3: Commit
git add cfc-backend/src/main/java/com/etotem/cfc/controller/ai/RecommendationController.java
git commit -m "feat(backend): 5% recommend traffic routed to LangGraph"
Files:
Create: cfc-langgraph/docker-compose.yml
[x] Step 1: 创建 docker-compose.yml
version: "3.8"
services:
langgraph-svc:
build:
context: .
dockerfile: Dockerfile
ports:
- "9000:9000"
env_file:
- .env
environment:
- JAVA_BASE_URL=http://host.docker.internal:9082
- CHROMA_DB_PATH=/data/chroma_db
volumes:
- langgraph-data:/data
- ./app:/app/app # 开发模式: 热重载
command: uvicorn app.main:app --host 0.0.0.0 --port 9000 --reload
restart: unless-stopped
volumes:
langgraph-data:
[x] Step 2: Commit
git add cfc-langgraph/docker-compose.yml
git commit -m "chore(langgraph): Docker Compose for local dev"
Files:
cfc-langgraph/tests/__init__.pycfc-langgraph/tests/conftest.pyCreate: cfc-langgraph/tests/test_recommend.py
[x] Step 1: 创建 tests/__init__.py(空文件)
[x] Step 2: 创建 tests/conftest.py
import pytest
from httpx import AsyncClient, ASGITransport
from app.main import app
@pytest.fixture
async def client():
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as ac:
yield ac
@pytest.fixture
async def health_response(client):
resp = await client.get("/health")
return resp
[x] Step 3: 创建 tests/test_recommend.py
import pytest
from httpx import AsyncClient, ASGITransport
from app.main import app
@pytest.mark.asyncio
async def test_health_check():
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
resp = await client.get("/health")
assert resp.status_code == 200
assert resp.json() == {"status": "ok"}
@pytest.mark.asyncio
async def test_recommend_empty_tags():
"""空标签应返回空列表"""
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
resp = await client.post("/api/v1/recommend", json={
"user_id": 1,
"tags": [],
"query": "推荐一些适合孩子的活动",
"limit": 3,
})
assert resp.status_code == 200
data = resp.json()
assert "items" in data
@pytest.mark.asyncio
async def test_recommend_with_tags():
"""带标签应返回推荐结果"""
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
resp = await client.post("/api/v1/recommend", json={
"user_id": 1,
"tags": ["益生菌", "肠胃"],
"limit": 3,
})
assert resp.status_code == 200
data = resp.json()
assert "items" in data
assert "source" in data
[x] Step 4: 运行测试
cd cfc-langgraph
pip install -e ".[dev]"
pytest tests/ -v
# 预期: 2 passed (test_health_check, test_recommend_empty_tags 通过)
# test_recommend_with_tags 可能因 Java 后端未运行而返回空列表, 但不应该报错
[x] Step 5: Commit
git add cfc-langgraph/tests/
git commit -m "test(langgraph): Phase 1 integration tests"