Researcher(공시 분석) → Analyst(기술적 분석) → Risk Manager(리스크 평가) → Portfolio Manager(최종 결정). 각 에이전트가 전문 역할을 맡아 LangGraph 워크플로우로 협업합니다.
포트폴리오 관리는 단일 에이전트로 처리하기 어려운 이유가 있습니다.
| 문제 | 단일 에이전트 한계 | 멀티 에이전트 해결 |
|---|---|---|
| 컨텍스트 길이 | 10개 종목 분석 시 컨텍스트 초과 | 에이전트별 독립 컨텍스트 |
| 전문화 | 공시 분석·기술 분석 동시에 잘 못함 | 에이전트별 역할 분리 |
| 병렬 처리 | 순차적으로만 처리 | 여러 종목 병렬 리서치 |
| 체크 & 밸런스 | 자기 오류 발견 어려움 | Risk Manager가 검토 |
설계 원칙: 멀티 에이전트는 꼭 필요한 경우에만 사용합니다. 이 시나리오(포트폴리오 리밸런싱)는 리서치·분석·리스크 평가·최종 결정이라는 명확한 역할 분리가 필요하므로 적합합니다.
[사용자 요청] "내 포트폴리오 리밸런싱 해줘"
│
▼
┌─────────────┐
│ Researcher │ DART/SEC 공시 분석, 뉴스 수집
└──────┬──────┘
│ research_report
▼
┌─────────────┐
│ Analyst │ RSI, MACD, 볼린저밴드 기술적 분석
└──────┬──────┘
│ analysis_report
▼
┌──────────────┐
│ Risk Manager │ 리스크 평가, 포지션 사이징
└──────┬───────┘
│ risk_report
▼
┌──────────────────┐
│ Portfolio Manager │ 최종 리밸런싱 플랜 생성
└──────────────────┘
│
▼
[최종 리포트 반환]from typing import TypedDict, Annotated
from langgraph.graph import add_messages
from langchain_core.messages import BaseMessage
class PortfolioState(TypedDict):
"""멀티 에이전트 간 공유 상태"""
# 사용자 입력
user_request: str
tickers: list[str] # 분석할 종목 목록
# 각 에이전트 출력
research_report: str # Researcher 결과
analysis_report: str # Analyst 결과
risk_report: str # Risk Manager 결과
final_plan: str # Portfolio Manager 최종 플랜
# 메시지 히스토리 (에이전트 간 전달)
messages: Annotated[list[BaseMessage], add_messages]from langchain_ollama import ChatOllama
from langchain_core.messages import HumanMessage
from rag.retriever import search_kr_disclosure, search_us_filing
from tools.us_stock import get_us_stock_price
from tools.kr_stock import get_kr_stock_price
from .state import PortfolioState
_researcher_llm = ChatOllama(model="llama3.1:8b", temperature=0)
_tools = [search_kr_disclosure, search_us_filing, get_us_stock_price, get_kr_stock_price]
_researcher_chain = _researcher_llm.bind_tools(_tools)
def researcher_node(state: PortfolioState) -> dict:
"""공시·뉴스 기반 펀더멘털 리서치"""
tickers = state["tickers"]
ticker_list = ", ".join(tickers)
prompt = f"""다음 종목들의 펀더멘털을 분석하세요: {ticker_list}
각 종목에 대해:
1. 최근 공시/실적 핵심 내용
2. 사업 현황 및 주요 변화
3. 투자 관점 정리 (긍정/부정 요인)
간결하게 종목별로 정리하세요."""
response = _researcher_chain.invoke([HumanMessage(content=prompt)])
# Tool 호출 결과 처리
tool_results = []
for tc in response.tool_calls if hasattr(response, "tool_calls") else []:
tool_fn = {t.name: t for t in _tools}.get(tc["name"])
if tool_fn:
result = tool_fn.invoke(tc["args"])
tool_results.append(f"[{tc['name']}] {result}")
research_content = response.content
if tool_results:
research_content += "\n\n[조회 데이터]\n" + "\n".join(tool_results)
return {"research_report": research_content}from langchain_ollama import ChatOllama
from langchain_core.messages import HumanMessage
from tools.technical import get_rsi, get_macd, get_bollinger_bands
from tools.signal import get_composite_signal
from .state import PortfolioState
_analyst_llm = ChatOllama(model="llama3.1:8b", temperature=0)
_tools = [get_composite_signal, get_rsi, get_macd, get_bollinger_bands]
_analyst_chain = _analyst_llm.bind_tools(_tools)
def analyst_node(state: PortfolioState) -> dict:
"""기술적 분석 — 진입 타이밍과 모멘텀 평가"""
tickers = state["tickers"]
research = state.get("research_report", "")
prompt = f"""다음 종목들의 기술적 분석을 수행하세요: {', '.join(tickers)}
각 종목별로 종합 시그널(get_composite_signal)을 먼저 확인한 뒤,
필요한 경우 RSI, MACD, 볼린저밴드를 추가 확인하세요.
펀더멘털 리서치 요약:
{research[:500]}
기술적 분석 결과를 종목별로 정리하고, 진입 타이밍에 대한 의견을 포함하세요."""
response = _analyst_chain.invoke([HumanMessage(content=prompt)])
tool_results = []
for tc in response.tool_calls if hasattr(response, "tool_calls") else []:
tool_fn = {t.name: t for t in _tools}.get(tc["name"])
if tool_fn:
result = tool_fn.invoke(tc["args"])
tool_results.append(f"[{tc['name']}] {result}")
analysis_content = response.content
if tool_results:
analysis_content += "\n\n[지표 데이터]\n" + "\n".join(tool_results)
return {"analysis_report": analysis_content}from langchain_ollama import ChatOllama
from langchain_anthropic import ChatAnthropic
from langchain_core.messages import HumanMessage
from .state import PortfolioState
import os
RISK_SYSTEM = """당신은 투자 리스크 관리 전문가입니다.
리서치 결과와 기술적 분석을 검토해 리스크를 평가합니다.
평가 항목:
- 종목별 최대 허용 비중 (변동성 기준)
- 섹터 집중 리스크
- 시장 상황 리스크 (금리, 환율, 지수 추세)
- 손절 기준 제안"""
def risk_manager_node(state: PortfolioState) -> dict:
"""리스크 평가 및 포지션 사이징 — Claude 사용 (복잡한 추론 필요)"""
# 리스크 평가는 Claude로 처리 (정교한 판단 필요)
llm = ChatAnthropic(
model="claude-sonnet-4-6",
api_key=os.getenv("ANTHROPIC_API_KEY"),
temperature=0,
max_tokens=2048,
)
prompt = f"""투자 리스크를 평가하고 포지션 사이징을 제안하세요.
분석 대상: {', '.join(state['tickers'])}
=== 펀더멘털 리서치 ===
{state.get('research_report', '')[:1000]}
=== 기술적 분석 ===
{state.get('analysis_report', '')[:1000]}
다음을 포함해 리스크 리포트를 작성하세요:
1. 종목별 리스크 등급 (상/중/하)
2. 권장 포트폴리오 비중 (총합 100%)
3. 핵심 리스크 요인
4. 손절 기준 제안"""
response = llm.invoke([HumanMessage(content=prompt)])
return {"risk_report": response.content}from langchain_anthropic import ChatAnthropic
from langchain_core.messages import HumanMessage
from .state import PortfolioState
import os
def portfolio_manager_node(state: PortfolioState) -> dict:
"""3개 에이전트 보고서를 종합해 최종 리밸런싱 플랜 생성"""
llm = ChatAnthropic(
model="claude-sonnet-4-6",
api_key=os.getenv("ANTHROPIC_API_KEY"),
temperature=0,
max_tokens=3000,
)
prompt = f"""사용자 요청: {state['user_request']}
아래 3개 전문가 보고서를 종합해 구체적인 포트폴리오 리밸런싱 플랜을 작성하세요.
=== 1. 펀더멘털 리서치 ===
{state.get('research_report', '')}
=== 2. 기술적 분석 ===
{state.get('analysis_report', '')}
=== 3. 리스크 평가 ===
{state.get('risk_report', '')}
최종 리밸런싱 플랜 형식:
## 포트폴리오 리밸런싱 플랜
### 종목별 액션
| 종목 | 현재 비중 | 목표 비중 | 액션 | 근거 |
|------|---------|---------|------|------|
### 우선순위 실행 순서
1. ...
### 리스크 관리
- 공통 손절 원칙: ...
### 주의사항
- 본 플랜은 참고용이며 투자 권유가 아닙니다"""
response = llm.invoke([HumanMessage(content=prompt)])
return {"final_plan": response.content}from langgraph.graph import StateGraph, END
from .state import PortfolioState
from .researcher import researcher_node
from .analyst import analyst_node
from .risk_manager import risk_manager_node
from .portfolio_manager import portfolio_manager_node
def build_portfolio_graph():
builder = StateGraph(PortfolioState)
# 노드 등록
builder.add_node("researcher", researcher_node)
builder.add_node("analyst", analyst_node)
builder.add_node("risk_manager", risk_manager_node)
builder.add_node("portfolio_manager", portfolio_manager_node)
# 엣지 (순차 실행)
builder.set_entry_point("researcher")
builder.add_edge("researcher", "analyst")
builder.add_edge("analyst", "risk_manager")
builder.add_edge("risk_manager", "portfolio_manager")
builder.add_edge("portfolio_manager", END)
return builder.compile()
# 싱글턴
_graph = None
def get_portfolio_graph():
global _graph
if _graph is None:
_graph = build_portfolio_graph()
return _graphfrom agents.portfolio.graph import get_portfolio_graph
graph = get_portfolio_graph()
# 포트폴리오 리밸런싱 요청
result = graph.invoke({
"user_request": "미국 AI 섹터와 한국 반도체 중심 포트폴리오를 리밸런싱해줘",
"tickers": ["NVDA", "AAPL", "MSFT", "005930", "000660"],
"messages": [],
"research_report": "",
"analysis_report": "",
"risk_report": "",
"final_plan": "",
})
print("=== 최종 리밸런싱 플랜 ===")
print(result["final_plan"])실행 시간: 5개 종목 기준 전체 파이프라인이 약 3~5분 소요됩니다. Researcher·Analyst는 Ollama(로컬), Risk Manager·Portfolio Manager는 Claude API를 사용해 품질과 속도를 균형 있게 맞췄습니다.
from fastapi import APIRouter, BackgroundTasks
from pydantic import BaseModel
from agents.portfolio.graph import get_portfolio_graph
import asyncio, uuid
router = APIRouter(prefix="/portfolio", tags=["포트폴리오"])
# 비동기 작업 상태 저장
_tasks: dict[str, str] = {}
class PortfolioRequest(BaseModel):
request: str
tickers: list[str]
class TaskResponse(BaseModel):
task_id: str
status: str
@router.post("/rebalance", response_model=TaskResponse)
async def request_rebalance(req: PortfolioRequest, bg: BackgroundTasks):
"""포트폴리오 리밸런싱 요청 (비동기 — 3~5분 소요)"""
task_id = str(uuid.uuid4())[:8]
_tasks[task_id] = "running"
def run():
graph = get_portfolio_graph()
result = graph.invoke({
"user_request": req.request,
"tickers": req.tickers,
"messages": [],
"research_report": "",
"analysis_report": "",
"risk_report": "",
"final_plan": "",
})
_tasks[task_id] = result["final_plan"]
bg.add_task(run)
return TaskResponse(task_id=task_id, status="running")
@router.get("/rebalance/{task_id}")
async def get_result(task_id: str):
"""리밸런싱 결과 폴링"""
status = _tasks.get(task_id, "not_found")
if status == "running":
return {"status": "running", "plan": None}
return {"status": "done", "plan": status}from agents.portfolio.graph import get_portfolio_graph
graph = get_portfolio_graph()
# 시나리오 1: AI 섹터 집중 포트폴리오
result = graph.invoke({
"user_request": "AI 반도체 중심 포트폴리오 리밸런싱",
"tickers": ["NVDA", "AMD", "005930", "000660"],
"messages": [], "research_report": "",
"analysis_report": "", "risk_report": "", "final_plan": "",
})
print(result["final_plan"])
# 시나리오 2: 한미 혼합 포트폴리오
result2 = graph.invoke({
"user_request": "달러 약세 환경에서 포트폴리오 방어 전략",
"tickers": ["AAPL", "MSFT", "005930", "035420", "207940"], # 삼성전자, NAVER, 삼성바이오
"messages": [], "research_report": "",
"analysis_report": "", "risk_report": "", "final_plan": "",
})
print(result2["final_plan"])# API로 비동기 실행
curl -X POST http://localhost:8000/portfolio/rebalance \
-H "Content-Type: application/json" \
-d '{
"request": "AI 섹터 집중 포트폴리오 리밸런싱",
"tickers": ["NVDA", "AAPL", "005930", "000660"]
}'
# {"task_id": "a1b2c3d4", "status": "running"}
# 3~5분 후 결과 조회
curl http://localhost:8000/portfolio/rebalance/a1b2c3d4