Automatyzacja wprowadzania transakcji w FastAPI przy użyciu asystenta AI LangGraph
Naucz się, jak stworzyć asystenta AI zarządzanego przez LangGraph, który analizuje język naturalny na strukturyzowane transakcje i zapisuje je do bazy danych PostgreSQL za pomocą FastAPI.
Każdy, kto próbował śledzić codzienne wydatki za pomocą formularza internetowego, wie, jak męczące to może być. Rejestracja zakupu kawy za 5 dolarów nie powinna wymagać wypełniania wielu pól, a jednak dokładnie to się dzieje w wielu aplikacjach do śledzenia finansów stworzonych z użyciem FastAPI i PostgreSQL, gdzie każda transakcja — bez względu na to, jak mała — musi być wprowadzana ręcznie.
Załóżmy, że budujemy taką aplikację, w której każda transakcja, zarówno prosta, jak i skomplikowana, musi przechodzić przez formularz. To działa dobrze przy sporadycznym wprowadzaniu danych, ale staje się wyczerpujące, gdy trzeba zarejestrować kilka transakcji za jednym razem.
Wtedy pojawia się naturalne pytanie: a co, gdyby cały ten proces mógł zostać zautomatyzowany? A co, gdyby zamiast wypełniać formularz, można było po prostu powiedzieć asystentowi „Dziś wydałem 5 dolarów na kawę w restauracji” i pozwolić mu zająć się resztą?
Tutaj właśnie przydaje się LangGraph.
Używając LangGraph, możesz stworzyć asystenta AI, który przyjmuje opis zwykłym językiem tego, co się wydarzyło, i przekształca go w prawidłowo zapisaną transakcję w twoim imieniu.
Wprowadzenie
W tym artykule omówiono podstawowe koncepcje leżące u podstaw LangGraph oraz otaczającego go ekosystemu, a następnie szczegółowo wyjaśniono, jak można dodać asystenta AI do aplikacji FastAPI przy użyciu LangGraph w celu automatyzacji wprowadzania transakcji.
LangGraph
LangGraph pochodzi od zespołu stojącego za LangChain i jest narzędziem otwartego oprogramowania służącym do budowy i zarządzania pracami agentów AI za pomocą struktur grafowych. Dzięki niemu proces można opisać jako zbiór „węzłów” i „krawędzi”, co sprawia, że złożone zachowania agentów są uporządkowane, skalowalne i łatwiejsze do kontrolowania.
Zanim przejdziemy dalej do LangGraph, przydatne jest najpierw zrozumienie LangChain, ponieważ LangGraph opiera się na nim.
LangChain
LangChain to również narzędzie otwartego oprogramowania służące do tworzenia aplikacji wykorzystujących duże modele językowe. Jego głównym zadaniem jest stworzenie mostu pomiędzy modelem językowym a zewnętrznymi zasobami — źródłami danych, narzędziami oraz krokami w procesie pracy — aby system mógł wykonywać rozumowanie wieloetapowe i zadania automatyczne, a nie ograniczać się do pojedynczej wymiany informacji w formie promptu i odpowiedzi.
Cel: Jest przeznaczony do tworzenia aplikacji AI, które wymagają łączenia kilku kroków — na przykład obsługi wprowadzonych przez użytkownika danych, pozyskiwania istotnych informacji oraz generowania odpowiedzi.
Struktura: LangChain opiera się na „łańcuchach”, które to są uporządkowane sekwencje operacji, w których wynik każdego kroku służy jako dane wejściowe dla następnego kroku. Dzięki temu można rozbić złożoną logikę na mniejsze, łatwiej zarządzalne części.
Zastosowania: Typowe przypadki użycia obejmują chatboty, zadania wymagające rozumowania wieloetapowego, wyszukiwanie i streszczanie dokumentów oraz łączenie modeli językowych z zewnętrznymi narzędziami lub API.
LangGraph (ciąg dalszy)
Mówiąc prościej, LangGraph organizuje wywołania modeli językowych w formie prac przepływowych o strukturze grafu, co umożliwia elastyczne, a nawet równoległe rozumowanie wieloetapowe zamiast ściśle liniowej sekwencji.
Cel: Umożliwia tworzenie aplikacji AI, w których logika może się rozgałęziać, tworzyć pętle lub wykonywać kroki równolegle, przekraczając to, co może wyrazić zwykły sekwencyjny łańcuch.
Struktura: LangGraph reprezentuje operacje jako „węzły”, a przepływ danych pomiędzy nimi jako „krawędzie”. Wyjście pojedynczego węzła może dostarczać dane do kilku kolejnych węzłów, co umożliwia dynamiczne ścieżki decyzyjne.
Zastosowania: LangGraph doskonale nadaje się do koordynacji wielu agentów, budowy złożonych procesów decyzyjnych, automatyzacji zadań wymagających logiki warunkowej oraz jednoczesnego zarządzania kilkoma modelami językowymi lub narzędziami.
Czym jest graf w LangGraph?
Graf to ogólnie struktura danych nieliniowa składająca się z „wierzchołków” (węzłów) i „krawędzi” (połączeń między nimi), która odzwierciedla relacje pomiędzy obiektami.
W LangGraph ta struktura grafu jest wykorzystywana do tworzenia stanowych, cyklicznych procesów pracy — takich, w których sztuczna inteligencja może podejmować decyzje, wracać do wcześniejszych kroków lub zmieniać kierunek w zależności od pośrednich wyników.
LangChain vs. LangGraph
Architektura projektu
(Koniec punktu dostępu FastAPI + LangGraph + tworzenie i przechowywanie transakcji)
Problem
Zanim LangGraph został wdrożony w aplikacji finansowej, tworzenie transakcji oznaczało bezpośrednie wywołanie punktu dostępu /transactions/add z treścią przypominającą tę:
{
title: "Coffee for Rosy",
type: "expense",
amount: 5,
note: "Paid $5 to Rosy for coffee",
category_id: "5bc22126-5982-4500-9e74-71c9c089f0c8",
payment_option_id: "07c5d180-fa4d-4435-aa04-b54ef436eca1"
}
Aby do tego dotrzeć, konieczne były dwa wcześniejsze wywołania API — jedno w celu pobrania listy kategorii, a drugie w celu pobrania opcji płatności — tylko po to, by uzyskać ID potrzebne do przesyłanych danych. Innymi słowy, utworzenie pojedynczej transakcji wymagało trzech kroków, przy czym były to kroki powolne.
Rozwiązanie
Rozwiązaniem było powierzenie wszystkich trzech kroków asystentowi AI, podczas gdy użytkownik musi jedynie opisać zwykłym językiem to, co zrobił ze swoimi pieniędzmi. Z myślą o tym celu poniżej przedstawiono strukturę implementacji.
Architektura trzech kroków
- Endpunkt FastAPI
- Orchestrator LangGraph
- Baza danych PostgreSQL
1. Endpunkt FastAPI
Użytkownik wysyła żądanie do punktu końcowego FastAPI /assistance/transaction-entry z treścią, która zawiera wiadomość opisującą transakcję.
{
message: "Sent $5 to Rosy for Coffee through cash."
}
2. Orchestrator LangGraph
Orchestrator jest zbudowany jako graf, w którym każdy węzeł reprezentuje operację, a każda krawędź przedstawia przepływ danych pomiędzy operacjami.
Pierwszy węzeł, Analyzer LLM, otrzymuje wiadomość użytkownika i sprawdza, czy zawiera ona wszystko niezbędne do zapisania transakcji — czy chodzi o przychód, czy wydatek, jaką jest kwota, jaki jest cel transakcji, jaka metoda płatności została użyta itp.
Jeśli wiadomość już zawiera wszystkie wymagane informacje, LLM przekształca ją na dane strukturyzowane, na przykład:
{
title: "Coffee",
transaction_type: "expense",
amount: 5.0,
note: "Sent $5 to Rosy for coffee through cash",
category: "coffee",
payment_option: "cash",
payment_type: "Cash",
is_complete: True, // Flag
missing_info_message: None // Flag
}
Takie ustrukturyzowane dane trafiają następnie do węzła Database Writer, po usunięciu dwóch pól flagowych (is_complete i missing_info_message). Węzeł Database Writer wywołuje metodę create_transaction(), która rejestruje transakcję związana z użytkownikiem w bazie danych.
Ale co się stanie, jeśli wiadomość nie zawiera pewnych szczegółów? Rozważmy taką wiadomość:
{
message: "Sent $5 to Rosy for Coffee." // payment mode is not specified
}
Tutaj tryb płatności nie jest określony. W tym przypadku dane wydobyte przez LLM będą zawierać uzupełnione wartości flag, takie jak:
{
title: "Coffee",
transaction_type: "expense",
amount: 5.0,
note: "Sent $5 to Rosy for coffee",
category: "coffee",
payment_option: None,
payment_type: None,
is_complete: False, // Flag
missing_info_message: "Please enter the payment mode used for this expense." // Flag
}
Ponieważ flaga is_complete ma tutaj wartość False, wiadomość missing_info_message jest kierowana do innego węzła połączonego z LLM Analyzer: węzła Clarification. Ta ścieżka jest aktywowana tylko wtedy, gdy is_complete ma wartość False.
Węzeł Clarification otrzymuje missing_info_message i wywołuje metodę ask_again(), która zwraca tę wiadomość jako odpowiedź na pierwotną prośbę FastAPI. Oznacza to koniec wykonywania grafu w tym przebiegu — użytkownik otrzymuje po prostu prośbę o uzupełnienie brakującej informacji, w tym przypadku sposobu płatności.
Załóżmy, że użytkownik odpowiada następnie brakującymi danymi, na przykład:
{
message: "UPI"
}
Taka odpowiedź powoduje ponowną inicjalizację grafu orkiestratora, który przechodzi przez tę samą sekwencję kroków co wcześniej.
Główna różnica w tym drugim etapie polega na tym, że żadne z wcześniejszych informacji nie ginie — historia rozmowy jest zachowywana za każdym razem, gdy dane są wyodrębniane przez LLM (ten mechanizm trwałości zostanie omówiony bardziej szczegółowo później w artykule). Ponieważ payment_option jest teraz dostępny i połączony z wcześniej zebranymi wartościami, is_complete zmienia się na True, a ostateczne, przefiltrowane dane są przekazywane do węzła Database Writer w następującej formie:
{
title: "Coffee",
transaction_type: "expense",
amount: 5.0,
note: "Sent $5 to Rosy for coffee through cash",
category: "coffee",
payment_option: "UPI",
payment_type: "Digital"
}
// Flags removed.
Następnie węzeł Database Writer przejmuje kontrolę, mając do dyspozycji te przefiltrowane dane. Należy pamiętać, że gdy wpis transakcji był tworzony ręcznie, wymagano dwóch dodatkowych wywołań API — jednego dla categories i drugiego dla payment_options — aby uzyskać odpowiednie ID przed utworzeniem samej transakcji. Ten sam problem występuje tutaj: przefiltrowane dane zawierają rzeczywiste wartości tekstowe kategorii i opcji płatności, a nie ich ID w bazie danych, a baza nie przyjmie surowych wartości dla tych pól.
Aby to rozwiązać, węzeł Database Writer musi zapytać bazę danych o odpowiednie rekordy kategorii i opcji płatności na podstawie wartości znajdujących się w przefiltrowanych danych.
Ponieważ projekt opiera się na FastAPI w połączeniu z SQLAlchemy, te zapytania są realizowane jako zapytania SQLAlchemy.
Dla categories:
from sqlalchemy import func, select
from sqlalchemy.exc import IntegrityError
# Run a select query to check if the category in data.category exists or not.
stmt = select(CategoriesModel).where(
CategoriesModel.user_id == user_id,
func.lower(getattr(CategoriesModel, name)) == data.category.lower()
)
# Execute the query.
result = await session.execute(stmt)
# If category exists, assign its ID to data.category.
existing_category = resule.scalar_one_or_none()
if existing_category:
data.category = existing_category.id
# If category doens't exits, create a new category and save it to database.
new_category = CategoriesModel(**{name: data.category, "user_id": user_id})
session.add(new_row)
try:
await session.flush()
except IntegrityError:
# In case another concurrent request created it first,
# we need to roll back and fetch it again.
await session.rollback()
result = await session.execute(stmt)
existing_category = result.scalar_one_or_none()
if existing_category:
data.category = existing_category.id
raise
Krótko mówiąc, ta logika:
Wykonuje zapytanie SELECT w celu sprawdzenia, czy kategoria wspomniana w data.category już istnieje.
Jeśli tak, ID tej kategorii zastępuje wartość w data.category.
Jeśli nie, tworzy się nowy rekord kategorii, a zamiast niego używany jest jej nowo wygenerowany ID.
To samo podejście stosuje się do payment_options:
from sqlalchemy import func, select
from sqlalchemy.exc import IntegrityError
# Run a select query to check if the peyment_option in data.payment_option exists or not.
stmt = select(PaymentOptionsModel).where(
PaymentOptionsModel.user_id == user_id,
func.lower(getattr(PaymentOptionsModel, name)) == data.payment_option.lower()
)
# Execute the query.
result = await session.execute(stmt)
# If payment_option exists, assign its ID to data.payment_option.
existing_option = resule.scalar_one_or_none()
if existing_option:
data.payment_option = existing_option.id
# If payment_option doesn't exits, create a new payment_option and save it to database.
new_option = PaymentOptionsModel(**{name: data.payment_option, "user_id": user_id})
session.add(new_row)
try:
await session.flush()
except IntegrityError:
# In case another concurrent request created it first,
# we need to roll back and fetch it again.
await session.rollback()
result = await session.execute(stmt)
existing_option = result.scalar_one_or_none()
if existing_option:
data.payment_option = existing_option.id
raise
Gdy już ustalono ID kategorii oraz ID opcji płatności, obiekt danych jest w pełni zaktualizowany i gotowy do wprowadzenia, wyglądając w ten sposób:
{
title: "Coffee",
type: "expense",
amount: 5.0,
note: "Sent $5 to Rosy for coffee through cash",
category_id: "5bc22126-5982-4500-9e74-71c9c089f0c8",
payment_option_id: "07c5d180-fa4d-4435-aa04-b54ef436eca1"
}
Z tymi ostatecznymi danymi węzeł Database Writer wywołuje metodę create_transaction(), która faktycznie zapisuje rekord transakcji w bazie danych.
3. Baza danych PostgreSQL
To przedstawia ostatni etap architektury, w którym sfinalizowane dane przekazane przez węzeł Database Writer są zapisywane do tabeli transactions.
Otrzymana struktura tabeli transactions wygląda następująco:
Implementacja
Gdy już omówiliśmy architekturę, nadszedł czas na przyjrzenie się szczegółom faktycznej implementacji tego orkiestratora przy użyciu LangGraph. Należy zauważyć, że kolejność przedstawiona tutaj nie odzwierciedla dokładnie powyższego opisu architektury. Zamiast tego implementacja jest zorganizowana w następujący sposób:
- LangGraph Orchestrator
- Baza danych PostgreSQL
- Endpunkt FastAPI
1. LangGraph Orchestrator
Sam orkiestrator znajduje się w pliku src/assistance/graph.py. Ten plik jest odpowiedzialny za konfigurację modelu LLM, definiowanie węzłów grafu, łączenie ich ze sobą oraz ostateczne skompilowanie całej struktury w działający graf.
Jak wspomniano wcześniej, orkiestrator składa się z trzech węzłów: Analyzer LLM, Database Writer oraz węzła Clarification.
Węzeł Analyzer LLM (Groq)
Węzeł ten jest w istocie modelem językowym, którego zadaniem jest zrozumienie intencji użytkownika oraz sprawdzenie, czy wiadomość zawiera wszystkie niezbędne i poprawne informacje. Zamiast budować własny model od zera, ten projekt polega na Groq do wykonywania tej trudnej pracy.
Czym jest Groq?
Groq to framework open-source napisany w języku Python, stworzony do pracy z danymi ustrukturyzowanymi w formie grafów. Daje programistom wyraźne narzędzia do wyszukiwania, filtrowania i agregacji informacji przechowywanych jako grafy, a doskonale nadaje się do obsługi dużych zbiorów danych grafowych, takich jak sieci społecznościowe, grafy wiedzy czy silniki rekomendacji, jak opisano w artykule GeekForGeeks na temat API Groq.
Dzięki hostowanemu API Groq można wysyłać zapytania do powszechnie używanych modeli otwartych — openai/gpt-oss-120b jest tym używanym w tym projekcie — i otrzymywać odpowiedzi, które zazwyczaj przychodzą znacznie szybciej niż te uzyskiwane od innych dostawców oferujących podobne modele.
Dlaczego Groq?
Groq został wybrany zamiast alternatyw takich jak ChatOpenAI czy ChatAnthropic z kilku powodów:
- Szybkość: Groq opiera się na specjalnie zaprojektowanej sprzętowo architekturze zwanej LPU (Language Processing Units), zamiast na GPU, od których zależą większość innych dostawców, co zapewnia szybkie przetwarzanie zapytań.
- Dostępna darmowa wersja: darmowa wersja Groq jest na tyle hojna, że pozwala realizować projekty indywidualne lub edukacyjne bez znaczących kosztów API podczas eksperymentowania.
- Zgodność dzięki LangChain: klasa
langchain_groq.ChatGroqłączy się z LangChain i LangGraph w taki sam sposób jakChatOpenAIlubChatAnthropic. Oznacza to, że przechodzenie na innego dostawcę później nie wymagałoby przepisywania logiki grafu — wystarczy zastąpić klienta.
Jak uzyskać klucz API Groq
Groq umożliwia generowanie darmowych kluczy API do celów rozwojowych. Oto jak je otrzymać:
- Wejdź na https://console.groq.com i zaloguj się lub utwórz konto.
- Wybierz opcję API Keys z paska nawigacyjnego.
- Wybierz opcję Create API Key.
- Pojawi się formularz z prośbą o podanie nazwy (w tym projekcie użyto
transaction-assistant) oraz okresu ważności klucza. Po wypełnieniu formularza należy go wysłać. - Klucz jest wyświetlany tylko raz, bezpośrednio po stworzeniu — dlatego koniecznie skopiuj go natychmiast.
Po utworzeniu wszystkie twoje klucze pojawią się na głównej liście na tej stronie.
Używanie klucza API Groq w kodzie FastAPI
Dodaj klucz API Groq do pliku .env znajdującego się w korzeniu projektu, obok innych zmiennych środowiskowych:
GROQ_API_KEY = "gsk_***************************************DyxM"
Istnieje kilka sposobów na ładowanie zmiennych środowiskowych do modułów, które je potrzebują. Ten projekt wykorzystuje dedykowaną klasę ustawień:
Zdefiniuj klasę Settings w pliku src/utils/settings.py:
from pydantic_settings import BaseSettings, SettingsConfigDict
class Settings(BaseSettings):
# Configure connection with the .env file
model_config = SettingsConfigDict(env_file=".env", extra="ignore")
# ... Other Variables ...
GROQ_API_KEY: str
settings = Settings()
Następnie importuj ten obiekt ustawień tam, gdzie jest on potrzebny:
from src.utils.settings import settings
# After importing, the object settings can be used as
# "settings.GROQ_API_KEY" to access the environment variable for Groq API Key.
Konfiguracja LLM
Zanim skonfigurujesz LLM, zainstaluj LangGraph i LangChain wraz z integracją Groq:
pip install -U langgraph langchain langchain-groq
Następnie tworzona jest instancja klienta Groq i konfiguruje się ją z określonym modelem:
from langchain_groq import ChatGroq
from src.assitance.schema import ExtractedTransactionSchema
from src.utils.settings import settings
assistance_llm = ChatGroq(model="openai/gpt-oss-120b", temperature=0.2,
api_key=settings.GROQ_API_KEY)
structured_llm = assistance_llm.with_structured_output(
ExtractedTransactionSchema)
Tutaj ChatGroq pełni rolę otulacza LangChain wokół modeli czatowych Groq, umożliwiając interakcję z nimi za pośrednictwem standardowej interfejsu LangChain zamiast ręcznego tworzenia żądań HTTP.
assistance_llm = ChatGroq(model="openai/gpt-oss-120b", temperature=0.2,
api_key=settings.GROQ_API_KEY)
Ten fragment tworzy wcześniej wspomnianą instancję klienta Groq, skonfigurowaną z wybranym modelem oraz wartością temperatury na niskim poziomie, a także uwierzytelnia się za pomocą klucza API pobranego z ustawień środowiskowych.
Temperatura to parametr, zazwyczaj w zakresie od 0 do 1, który określa, na ile losowe lub kreatywne są odpowiedzi modelu. Wyższa wartość, np. 0,8, skłania wyniki do większej różnorodności i kreatywności, natomiast niższa wartość, np. 0,2, sprawia, że odpowiedzi są bardziej spójne i przewidywalne. W tym projekcie ustawiono temperature = 0,2.
structured_llm = assistance_llm.with_structured_output(
ExtractedTransactionSchema)
Ten kod otacza model LLM w taki sposób, że zamiast zwracać zwykły tekst, generuje obiekt w języku Python, który dokładnie odpowiada ExtractedTransactionSchema. Wewnętrznie osiąga się to poprzez polecenie modelowi wygenerowania wyniku zgodnego ze schematem, a następnie automatyczną analizę i weryfikację tego wyniku — co eliminuje konieczność ręcznej interpretacji surowego tekstu modelu.
Sam ExtractedTransactionSchema jest zdefiniowany w pliku src/assistance/schema.py:
from typing import Optional
from pydantic import BaseModel, Field
class ExtractedTransactionSchema(BaseModel):
is_complete: bool = Field(
description="True only if title, type, amount, category, and payment method were all found.")
missing_info_message: Optional[str] = Field(
default=None, description="A polite clarifying question listing listing exactly what's missing. Must be null if is_complete is True")
title: str
transaction_type: str = Field(description="'income' or 'expense'")
amount: float
category: str
payment_option: str = Field(
description="e.g. 'UPI', 'Cash', 'HDFC Credit Card'")
payment_type: str = Field(
description="Broad classification of the payment_option, one of: 'Cash', 'Card', 'Digital', 'Bank Transfer', 'Other'"
)
note: str
Należy zauważyć, że w tym momencie model LLM jeszcze nie został wywołany — ten krok jedynie określa formę, jaką musi mieć wynik po jego uruchomieniu.
Stan grafu
Stan grafu reprezentuje strukturę danych, która przepływa przez graf i jest aktualizowana w jego obrębie. Można go porównać do pamięci roboczej dyrygenta: przechowuje on każdą informację, którą graf śledzi i modyfikuje w miarę postępów wykonywania poszczególnych kroków. Dla tego asystenta transakcyjnego stan grafu jest zdefiniowany w następujący sposób:
from pydantic import BaseModel, Field
from typing import Annotated, List, Optional
import operator
from src.assitance.schema import ExtractedTransactionSchema
class GraphState(BaseModel):
user_input: str = Field(description="The user input to the graph.")
conversation_history: Annotated[List[str], operator.add] = []
extracted: Optional[ExtractedTransactionSchema] = None
final_response: Optional[str] = None
Rozbierzmy, co tak naprawdę robi ten kod:
from pydantic import BaseModel, Field
Pydantic to biblioteka do walidacji danych używana tutaj. BaseModel to klasa nadrzędna, którą rozszerza się podczas definiowania ustrukturyzowanej, sprawdzanej typowo formy, takiej jak GraphState. Field umożliwia dołączanie metadanych — opisów, wartości domyślnych itp. — do każdego pojedynczego atrybutu.
from typing import Annotated, List, Optional
import operator
To są narzędzia do typowania w Pythonie. Optional wskazuje, że pole może być puste i zawierać wartość None. List określa atrybut jako listę elementów. Annotated, w połączeniu z operator.add, informuje LangGraph „gdy węzeł zwraca nową wartość dla tego pola, należy dodać ją do istniejącej listy zamiast ją zastępować”. To właśnie ten mechanizm umożliwia, aby conversation_history rosła z każdą kolejną rozmową, zamiast być wymazywana przy każdej nowej wiadomości.
class GraphState(BaseModel):
user_input: str = Field(description="The user input to the graph.")
conversation_history: Annotated[List[str], operator.add] = []
extracted: Optional[ExtractedTransactionSchema] = None
final_response: Optional[str] = None
user_input: najnowsza wiadomość wysłana przez użytkownika w ramach tej konkretnej interakcji.conversation_history: pełna lista wcześniejszych wiadomości, gromadzona krok po kroku zamiast była nadpisywana.
extracted: uzupełniane po tym, jak model językowy wydobyje ze rozmowy strukturyzowane dane transakcyjne. Początkowo ma wartość None, ponieważ na początku wykonywania nie zostało jeszcze nic wydobyte.final_response: wiadomość, która ostatecznie trafia z powrotem do użytkownika — albo potwierdzenie, że transakcja została zapisana, albo pytanie dodatkowe wymagające więcej szczegółów.Zapytanie do wydobywania danych
from langchain_core.prompts import PromptTemplate
EXTRACTION_PROMPT = PromptTemplate(
template="""
You are a financial assistant extracting transaction details.
Below is the conversation so far (it may span multiple messages, where later
messages answer questions raised by earlier ones). Treat it as one combined input.
Required fields: title, transaction_type (income/expense), amount, category, payment_option.
If title is missing, add one based on the context of the message.
If anything required is missing, except title, set is_complete to False and write a short, polite
clarifying question in missing_info_message asking only for what's missing.
If everything is present, set is_complete to True, leave missing_info_message null,
and fill in all fields. Always copy the user's original message into `note`.
Conversation so far:
{user_input}
""",
input_variables=["user_input"]
)
To jest dosłowna instrukcja przekazywana modelowi językowemu w języku naturalnym — określa, jakie pola należy sprawdzić, co robić w przypadku braku danych oraz jak powinna być zorganizowana odpowiedź. Ponieważ structured_llm już egzekwuje schemat na poziomie wyjścia, zadaniem promptu jest głównie kierowanie rozumowaniem modelu: decydowanie o tym, co oznacza „kompletność”, jak sformułować pytanie wyjaśniające itp., podczas gdy schemat zajmuje się formatowaniem.
Ekstraktor
def extractor(state: GraphState):
full_conversation = "\n".join(
state.conversation_history + [state.user_input])
prompt = EXTRACTION_PROMPT.format(user_input=full_conversation)
result: ExtractedTransactionSchema = structured_llm.invoke(prompt)
return {"extracted": result, "conversation_history": [state.user_input]}
Funkcja extractor wykonuje następujące czynności:
- Łączy wszystkie poprzednie wiadomości z bieżącą, aby model mógł zobaczyć pełny kontekst.
- Przekazuje ten połączony tekst modelowi językowemu.
- Otrzymuje z powrotem zstrukturyzowany obiekt
ExtractedTransactionSchema.
operator.add skonfigurowanemu dla tego pola.Decyzja
route_after_extraction
def route_after_extraction(state: GraphState):
return "create_transaction" if state.extracted.is_complete else "ask_again"
Ta funkcja nie wykonywa żadnej rzeczywistej obróbki — jej jedynym zadaniem jest podjęcie decyzji. W zależności od tego, czy LLM oznaczył wyekstrahowane dane jako kompletne, zwraca ciąg znaków wskazujący grafowi, który węzeł powinien zostać uruchomiony dalej. Można to porównać do logiki rozgałęzień w schemacie przepływu: graf sprawdza wartość zwróconą przez tę funkcję i podąża za odpowiadającą ścieżką, albo kieruje się do create_transaction, aby zarejestrować transakcję, albo do ask_again, aby poprosić o dodatkowe informacje.
Węzeł zapisujący do bazy danych
create_transaction_node
Ten węzeł jest odpowiedzialny za zapisywanie ukończonych danych transakcyjnych do bazy danych pod konto odpowiedniego użytkownika. Funkcja create_transaction_node, która implementuje węzeł DB Writer, wygląda następująco:
from langchain_core.runnables import RunnableConfig
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.ext.asyncio import AsyncSession
from src.transaction import controller
from src.transaction.schema import TransactionCreateSchema
from src.utils.db_helper import get_or_create
from src.categories.models import CategoriesModel
from src.categories.controller import get_deterministic_color
from src.payment_options.models import PaymentOptionsModel
from src.utils.db_helper import get_or_create
async def create_transaction_node(state: GraphState, config: RunnableConfig):
session: AsyncSession = config["configurable"]["session"]
user = config["configurable"]["user"]
data = state.extracted
try:
category_id = await get_or_create(
session, CategoriesModel, user.id, data.category,
extra_defaults={"color": get_deterministic_color(data.category)}
)
payment_option_id = await get_or_create(
session, PaymentOptionsModel, user.id, data.payment_option,
extra_defaults={"payment_type": data.payment_type}
)
payload = TransactionCreateSchema(
amount=data.amount,
category_id=category_id,
payment_option_id=payment_option_id,
note=data.note,
title=data.title,
type=data.transaction_type,
)
await controller.create_transaction(payload, session, user)
await session.commit()
except SQLAlchemyError as err:
await session.rollback()
print(
f"Error while creating transaction through AI assistance :: {err}")
return {
"final_response": "Something went wrong while saving your transaction. Please try again."
}
message = f"Added {data.transaction_type} of {data.amount} under '{data.category}' ({data.payment_option})"
return {"final_response": message}
To za dużo do przyswojenia naraz, więc przeanalizujmy to krok po kroku.
from langchain_core.runnables import RunnableConfig
Typ reprezentujący obiekt config przekazywany do dowolnego węzła. Istnieje on wyłącznie jako wskazówka typu, dzięki czemu każdy, kto czyta sygnaturę create_transaction_node, od razu rozumie, jaki kształt ma config.
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.ext.asyncio import AsyncSession
To standardowe importy z SQLAlchemy potrzebne do łapania błędów bazy danych oraz do określenia typu asynchronicznego sesji bazy danych używanej do komunikacji z PostgreSQL.
from src.transaction import controller
from src.transaction.schema import TransactionCreateSchema
To wykorzystuje logikę tworzenia transakcji, która jest już używana w innych częściach aplikacji, wraz z jej schematem danych wejściowych. Ponowne użycie tej logiki oznacza, że asystent tworzy transakcje za pomocą dokładnie tego samego ścieżki kodu co zwykła API CRUD, zamiast kopiować tę logikę tutaj.
from src.utils.db_helper import get_or_create
from src.categories.models import CategoriesModel
from src.categories.controller import get_deterministic_color
from src.payment_options.models import PaymentOptionsModel
Są to elementy pomocnicze służące do przekształcenia nazw kategorii i opcji płatności wyodrębnionych przez model językowy w rzeczywiste rekordy i ID w bazie danych, tworząc nowe wpisy, jeśli jeszcze ich nie ma.
Uwaga: logika wyszukiwania/tworzenia dla categories i payment_options została połączona w jeden uniwersalny pomocnik get_or_create, ponieważ oba modele wymagały w zasadzie tego samego zachowania.
async def create_transaction_node(state: GraphState, config: RunnableConfig):
Funkcja create_transaction_node jest wykonywana tylko wtedy, gdy potwierdzono, że dane zostały w pełni wydobyte. Jest zadeklarowana jako async, ponieważ wykonuje rzeczywiste operacje w bazie danych, a do jej działania potrzebne są config oraz state, aby mogła uzyskać dostęp do aktywnej sesji bazy danych i zalogowanego użytkownika. Te dwa wartości pochodzą z trasy API, a nie z LLM czy stanu rozmowy, ponieważ należą do konkretnego żądania, a nie do trwającej rozmowy.
session: AsyncSession = config["configurable"]["session"]
user = config["configurable"]["user"]
data = state.extracted
To wyodrębnia sesję, użytkownika oraz dane transakcyjne wydobyte z bazy danych.
try:
category_id = await get_or_create(...)
payment_option_id = await get_or_create(...)
Ponieważ LLM wyekstrahowało jedynie nazwy kategorii i metody płatności, takie jak „Artykuły spożywcze” lub „UPI”, a nie ich identyfikatory w bazie danych, ten krok sprawdza, czy dla bieżącego użytkownika już istnieje odpowiadająca wiersz. Jeśli nie, tworzy go. W obu przypadkach zwraca odpowiadający identyfikator.
payload = TransactionCreateSchema(...)
await controller.create_transaction(payload, session, user)
await session.commit()
Ładunek jest konstruowany w tym samym formacie, jakiego oczekuje istniejąca logika tworzenia transakcji, a następnie przekazywany do tej samej funkcji kontrolera, dzięki czemu wykorzystuje się istniejącą logikę aplikacji zamiast ją przepisywać. Następnie transakcja bazodanowa jest zatwierdzana, aby utrwalić zmianę.
except SQLAlchemyError as err:
await session.rollback()
print(...)
return {"final_response": "Something went wrong..."}
Jeśli coś pójdzie nie tak na warstwie bazy danych, wszystkie częściowe zmiany są cofane, a zamiast dopuścić do awarii żądania, zwracana jest przyjazna wiadomość o błędzie. Dzięki temu unika się sytuacji, w której na przykład powstaje nowa kategoria, ale nie ma do niej odpowiadającej transakcji.
message = f"Added {data.transaction_type} of {data.amount} under '{data.category}' ({data.payment_option})"
return {"final_response": message}
W przypadku sukcesu tworzona jest wiadomość potwierdzenia czytelna dla człowieka, która jest zwracana jako aktualizacja stanu.
Węzeł wyjaśniający
ask_again
def ask_again_node(state: GraphState):
return {"final_response": state.extracted.missing_info_message}
To jest prosty plan awaryjny. Jak opisano wcześniej, ten węzeł uruchamia się tylko wtedy, gdy wydobyte dane mają ustawioną wartość is_complete na False, wraz z przydatną wiadomością zawartą w missing_info_message.
Wewnątrz ask_again_node funkcja otrzymuje parametr state, dzięki czemu ma dostęp do state.extracted.is_complete oraz state.extracted.missing_info_message.
Krótko mówiąc, gdy brakuje informacji, ten węzeł po prostu przekazuje pytanie uzupełniające, które model językowy już wygenerował podczas extrakcji, aby użytkownik dokładnie wiedział, co należy dostarczyć dalej.
Budowanie grafu
from langgraph.graph import StateGraph, START, END
graph_builder = StateGraph(GraphState)
To tworzy nowego budowniczego grafu i informuje go, że każdy węzeł w grafie będzie czytał z obiektu o strukturze GraphState oraz do niego zapisywał dane.
graph_builder.add_node("extractor", extractor)
graph_builder.add_node("create_transaction", create_transaction_node)
graph_builder.add_node("ask_again", ask_again_node)
Każda funkcja jest tutaj rejestrowana jako nazwany węzeł, czyli zasadniczo oznaczony krok w obrębie grafu.
graph_builder.add_edge(START, "extractor")
To określa punkt wejścia: każda eksploatacja grafu rozpoczyna się od węzła extractor.
graph_builder.add_conditional_edges(
"extractor",
route_after_extraction,
{
"create_transaction": "create_transaction",
"ask_again": "ask_again",
},
)
Tutaj odbywa się rozgałęzianie. Gdy extractor zakończy swoją pracę, LangGraph wywołuje route_after_extraction, aby określić następny krok. Niezależnie od tego, jaką wartość zwraca – create_transaction czy ask_again – jest ona sprawdzana w tym zestawieniu, które łączy każdą wartość decyzyjną z konkretnym węzłem, do którego należy przejść.
graph_builder.add_edge("create_transaction", END)
graph_builder.add_edge("ask_again", END)
Oba możliwe rozgałęzienia zakończą wykonywanie grafu po swoim ukończeniu, przy czym w obu przypadkach zostanie osiągnięty punkt END.
Kompilowanie z wykorzystaniem pamięci
from langgraph.checkpoint.memory import MemorySaver
memory = MemorySaver()
assistance_graph = graph_builder.compile(checkpointer=memory)
Wywołanie compile() przekształca definicję grafu w coś, co można uruchomić. Przekazanie parametru checkpointer=memory aktywuje opisany wcześniej mechanizm utrzymywania stanu, dzięki czemu ponowne wywołanie grafu z tym samym thread_id kontynuuje pracę od miejsca przerwania, zamiast ją restartować.
Ostateczny kod
Tym samym warstwa orkiestracji jest gotowa. Oto finalny plik (src/assistance/graph.py):
from langgraph.graph import StateGraph, START, END
from langchain_groq import ChatGroq
from langgraph.checkpoint.memory import MemorySaver
from langchain_core.prompts import PromptTemplate
from langchain_core.runnables import RunnableConfig
from pydantic import BaseModel, Field
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.ext.asyncio import AsyncSession
from typing import Annotated, List, Optional
import operator
from src.utils.settings import settings
from src.assitance.schema import ExtractedTransactionSchema
from src.transaction import controller
from src.transaction.schema import TransactionCreateSchema
from src.utils.db_helper import get_or_create
from src.categories.models import CategoriesModel
from src.categories.controller import get_deterministic_color
from src.payment_options.models import PaymentOptionsModel
assistance_llm = ChatGroq(model="openai/gpt-oss-120b", temperature=0.2,
api_key=settings.GROQ_API_KEY)
structured_llm = assistance_llm.with_structured_output(
ExtractedTransactionSchema)
class GraphState(BaseModel):
user_input: str = Field(description="The user input to the graph.")
conversation_history: Annotated[List[str], operator.add] = []
extracted: Optional[ExtractedTransactionSchema] = None
final_response: Optional[str] = None
EXTRACTION_PROMPT = PromptTemplate(
template="""
You are a financial assistant extracting transaction details.
Below is the conversation so far (it may span multiple messages, where later
messages answer questions raised by earlier ones). Treat it as one combined input.
Required fields: title, transaction_type (income/expense), amount, category, payment_option.
If title is missing, add one based on the context of the message.
If anything required is missing, except title, set is_complete to False and write a short, polite
clarifying question in missing_info_message asking only for what's missing.
If everything is present, set is_complete to True, leave missing_info_message null,
and fill in all fields. Always copy the user's original message into `note`.
Conversation so far:
{user_input}
""",
input_variables=["user_input"]
)
def extractor(state: GraphState):
full_conversation = "\n".join(
state.conversation_history + [state.user_input])
prompt = EXTRACTION_PROMPT.format(user_input=full_conversation)
result: ExtractedTransactionSchema = structured_llm.invoke(prompt)
return {"extracted": result, "conversation_history": [state.user_input]}
def route_after_extraction(state: GraphState):
return "create_transaction" if state.extracted.is_complete else "ask_again"
async def create_transaction_node(state: GraphState, config: RunnableConfig):
session: AsyncSession = config["configurable"]["session"]
user = config["configurable"]["user"]
data = state.extracted
try:
category_id = await get_or_create(
session, CategoriesModel, user.id, data.category,
extra_defaults={"color": get_deterministic_color(data.category)}
)
payment_option_id = await get_or_create(
session, PaymentOptionsModel, user.id, data.payment_option,
extra_defaults={"payment_type": data.payment_type}
)
payload = TransactionCreateSchema(
amount=data.amount,
category_id=category_id,
payment_option_id=payment_option_id,
note=data.note,
title=data.title,
type=data.transaction_type,
)
await controller.create_transaction(payload, session, user)
await session.commit()
except SQLAlchemyError as err:
await session.rollback()
return {
"final_response": "Something went wrong while saving your transaction. Please try again."
}
message = f"Added {data.transaction_type} of {data.amount} under '{data.category}' ({data.payment_option})"
return {"final_response": message}
def ask_again_node(state: GraphState):
return {"final_response": state.extracted.missing_info_message}
graph_builder = StateGraph(GraphState)
graph_builder.add_node("extractor", extractor)
graph_builder.add_node("create_transaction", create_transaction_node)
graph_builder.add_node("ask_again", ask_again_node)
graph_builder.add_edge(START, "extractor")
graph_builder.add_conditional_edges(
"extractor",
route_after_extraction,
{
"create_transaction": "create_transaction",
"ask_again": "ask_again",
},
)
graph_builder.add_edge("create_transaction", END)
graph_builder.add_edge("ask_again", END)
memory = MemorySaver()
assistance_graph = graph_builder.compile(checkpointer=memory)
2. Baza danych PostgreSQL
Interakcja z bazą danych na tym etapie jest już obsługiwana wewnątrz funkcji create_transaction_node, omówionej wcześniej, gdzie sfinalizowane dane transakcyjne są zapisywane do tabeli transactions.
3. Koniec punktu końcowego FastAPI
from fastapi import APIRouter, Depends, status
from sqlalchemy.ext.asyncio import AsyncSession
from src.assitance.schema import UserMessageSchema
from src.assitance.graph import assistance_graph
from src.auth.models import UsersModel
from src.utils.db import get_db
from src.utils.auth.authentication import allow_all
assistance_routes = APIRouter(prefix="/assistance")
@assistance_routes.post("/transaction-entry", status_code=status.HTTP_201_CREATED)
async def run_transaction_assistance(payload: UserMessageSchema, session: AsyncSession = Depends(get_db), user: UsersModel = Depends(allow_all)):
config = {"configurable": {
"thread_id": str(user.id),
"session": session,
"user": user
}
}
result = await assistance_graph.ainvoke(
{"user_input": payload.message}, config=config)
return {"response": result["final_response"]}
Klienci wywołują ten endpoint (/assistance/transaction-entry) i do ciała żądania dodają wiadomość opisującą, o co chodziło w transakcji.
Rozważmy, jakie funkcje pełni każda z tych części.
@assistance_routes.post("/transaction-entry", status_code=status.HTTP_201_CREATED)
To konfiguruje ścieżkę POST pod adresem /transaction-entry. Ustawienie status_code=status.HTTP_201_CREATED mówi FastAPI, jaki kod stanu ma być domyślnie zwracany w przypadku sukcesu. 201 to standardowy kod oznaczający „stworzono nowy zasób”, co pasuje tutaj, ponieważ udane wywołanie skutkuje utworzeniem nowej wiersza transakcji.
Konstruowanie konfiguracji grafu:
config = {
"configurable": {
"thread_id": str(user.id),
"session": session,
"user": user
}
}
To tworzy obiekt config, który jest przekazywany podczas wywołania grafu. Obiekt config zawiera wartości specyficzne dla żądania, które nie powinny znajdować się w samym utrwalonym stanie rozmowy.
"thread_id": str(user.id): ta wartość jest wykorzystywana przez mechanizm checkpointer w LangGraph do określenia, czy historię rozmowy którego użytkownika należy pobrać i zaktualizować. Dzięki powiązaniu jej z identyfikatorem zalogowanego użytkownika każdy użytkownik automatycznie otrzymuje izolowany, trwały wątek rozmowy, dzięki czemu nieukończona transakcja jednego użytkownika nigdy nie może wpłynąć na innego. Wartość ta jest przekształcana w ciąg znaków, ponieważ checkpointer oczekujethread_idw formie ciągu znaków, podczas gdyuser.idzazwyczaj ma postać UUID."session"i"user": są one przekazywane dalej, aby funkcjacreate_transaction_node, działająca wewnątrz grafu, miała dostęp do aktualnej sesji bazy danych oraz do informacji o tym, kto wysyła żądanie.
Wywoływanie grafu:
result = await assistance_graph.ainvoke(
{"user_input": payload.message}, config=config)
To jest wiersz, który faktycznie uruchamia wykonywanie. ainvoke jest asynchronicznym odpowiednikiem uruchamiania grafu, natomiast użycie synchronicznego invoke spowodowałoby zablokowanie pętli zdarzeń, co jest istotne w tym przypadku, ponieważ create_transaction_node wewnętrznie wykonywa operacje asynchroniczne z bazą danych.
- Pierwszy argument,
{"user_input": payload.message}, reprezentuje początkowy stan tej sesji. Wyraźnie należy podać tylkouser_input; pozostałe pola typuGraphState(conversation_history,extracted,final_response) mają albo wartości domyślne, albo są wypełniane w miarę przetwarzania danych przez graf. Gdy istniejącythread_idjuż zawiera zapisaną historię, LangGraph łączy ten nowy wprowadzony danych z tym zapisanym stanem, zamiast rozpoczynać od zera. config=configdostarcza wszystko, co zostało przygotowane na poprzednim kroku:thread_idsłużący do znalezienia odpowiedniego stanu, a takżesessioniuserdla węzła odpowiedzialnego za zapis do bazy danych.
await wstrzymuje tę korekturę do chwili zakończenia pełnego wykonywania grafu, ponieważ ainvoke zwraca korekturę, której należy oczekiwać przed użyciem wyniku.To, co zostaje zwrócone jako result, to ostateczny stan GraphState, reprezentowany jako słownik, odzwierciedlający wykonywanie grafu, niezależnie od tego, czy zakończyło się ono w funkcji create_transaction, czy w funkcji ask_again.
Zwracanie odpowiedzi:
return {"response": result["final_response"]}
Trasa kończy się zwróceniem prostego słownika zawierającego jedynie tekst ostatecznej odpowiedzi. FastAPI zajmuje się konwersją tego tekstu na payload JSON dla klienta, tworząc coś w rodzaju:
{ "response": "Added expense of 450 under 'Groceries' (UPI)" }
Ten tekst jest identyczny z tym, który został wcześniej utworzony w funkcjach create_transaction_node lub ask_again_node. Sama ścieżka nie wie, która z gałęzi została faktycznie wykorzystana; po prostu przekazuje to, co trafiło do zmiennej final_response.
Wniosek
Praca nad tym asystentem transakcyjnym uwydatnia coś, co w materiałach instruktażowych często pomija się: trudną częścią wdrożenia funkcji opartej na AI nie jest proszenie modelu o odpowiedź, lecz zapewnienie, że ta odpowiedź będzie bezpiecznie funkcjonować po interakcji z rzeczywistym systemem. Tworzenie pytań jest łatwe. Prawdziwy wysiłek inżynierski wymaga opracowania schematów gwarantujących ustrukturyzowany wynik, grafów warunkowych decydujących o tym, czy zapisać dane, czy poprosić o dodatkowe wyjaśnienia, oraz stanu, który jest prawidłowo przenoszony między kolejnymi rundami rozmowy.
LangGraph okazał się odpowiedni dla tego projektu właśnie dlatego, że proces pracy wymagał rzeczywistego podejmowania decyzji, a nie tylko jednokierunkowego przetwarzania danych od wejścia do wyjścia. Jeśli twoja funkcja wymaga prostego przepływu, zwykłe wywołanie modelu językowego lub łańcuch LangChain jest prawdopodobnie prostszym i bardziej odpowiednim narzędziem. Jednak gdy logika AI musi się rozgałęziać, przechowywać informacje lub zawiesić pracę w celu zebrania dodatkowych danych przed kontynuacją, struktura oparta na grafach przestaje wyglądać jak zbędna złożoność i staje się najrozsądniejszym sposobem na odwzorowanie tego przepływu.
Odnośniki
Pozycje pokrewne
- LangChain vs LangGraph: Wybór między łańcuchami a grafami z stanem — Dowiedz się, w jaki sposób liniowe elementy budulcowe LangChain różnią się od zorientowanych na stan i rozgałęzionych procesów pracy LangGraph, oraz jak zdecydować, który z nich pasuje do twojej aplikacji AI.
- Budowanie agenta AI od zera: wzory, ReAct i LangGraph — Poznaj podstawowe koncepcje agentów AI – planowanie, używanie narzędzi, refleksja oraz wzorzec ReAct – oraz to, jak LangChain i LangGraph mogą pomóc przy ręcznym tworzeniu takiego agenta.