Автоматизация ввода транзакций в FastAPI с использованием ИИ-ассистента LangGraph
Изучите, как создать вспомогательного ИИ-ассистента на базе LangGraph, который будет парсить естественный язык в структурированные транзакции и записывать их в базу данных PostgreSQL с помощью FastAPI.
Каждый, кто пробовал отслеживать ежедневные расходы с помощью веб-формы, знает, насколько это утомительно. Для записи покупки кофе стоимостью 5 долларов не должно требоваться заполнение множества полей, однако именно это происходит во многих приложениях для отслеживания финансов, созданных с использованием FastAPI и PostgreSQL, где каждая транзакция — какой бы незначительной она ни была — должна вводиться вручную.
Представьте себе создание такого приложения, где каждая транзакция, будь она простой или сложной, должна проходить через форму. Это подходит для единичных записей, но становится изнурительным, когда нужно зафиксировать несколько транзакций сразу.
В такой ситуации возникает естественный вопрос: а что, если весь процесс можно автоматизировать? Что, если вместо заполнения формы можно просто сказать ассистенту: «Сегодня я потратил 5 долларов на кофе в ресторане», и позволить ему заниматься остальным?
Именно здесь и пригождается LangGraph.
С помощью LangGraph вы можете создать ИИ-ассистента, который принимает описание произошедшего на простом языке и преобразует его в соответствующую зарегистрированную транзакцию от вашего имени.
Введение
В этой статье рассматриваются основные концепции LangGraph и связанной с ним экосистемы, а затем подробно объясняется, как можно добавить ИИ-ассистента в приложение FastAPI с использованием LangGraph для автоматизации ввода транзакций.
LangGraph
LangGraph был создан командой, разработавшей LangChain, и представляет собой инструментарий с открытым исходным кодом для формирования и управления рабочими процессами ИИ-агентов с помощью графовых структур. С его помощью процесс описывается как совокупность «узлов» и «ребер», что позволяет сохранять сложное поведение агентов организованным, масштабируемым и удобным в управлении.
Прежде чем углубляться в LangGraph, полезно сначала понять LangChain, поскольку LangGraph строится на его основе.
LangChain
LangChain — это также инструментарий с открытым исходным кодом для создания приложений, работающих с большими языковыми моделями. Его основная функция — служить мостом для разработчиков между LLM и внешними ресурсами: источниками данных, инструментами и этапами рабочего процесса, что позволяет системе выполнять многоэтапные логические операции и автоматизированные задачи вместо простого обмена сообщениями в ответ на запрос.
Цель: Он предназначен для создания ИИ-приложений, которым необходимо объединять несколько этапов — например, обработка ввода пользователя, поиск соответствующей информации и формирование ответа.
Структура: LangChain основан на «цепях», представляющих собой упорядоченные последовательности операций, при которых результат каждого шага становится входными данными для следующего. Это позволяет разбивать сложную логику на более мелкие, управляемые части.
Применение: Типичные сценарии использования включают чат-ботов, задачи многошагового рассуждения, поиск и краткое изложение документов, а также подключение больших языковых моделей к внешним инструментам или API.
LangGraph (продолжение)
Проще говоря, LangGraph организует вызовы больших языковых моделей в рабочие процессы в форме графа, что обеспечивает гибкое и даже параллельное многошаговое рассуждение вместо строго линейной последовательности.
Цель: Он позволяет создавать ИИ-приложения, в которых логика может разветвляться, формировать циклы или выполнять шаги параллельно, выходя за рамки того, что может передать простая последовательная цепь.
Структура: LangGraph представляет операции в виде «узлов», а поток данных между ними — в виде «рёбер». Выходной сигнал одного узла может поступать в несколько следующих за ним узлов, что позволяет создавать динамические пути принятия решений.
Применения: LangGraph идеально подходит для координации нескольких агентов, создания сложных цепочек принятия решений, автоматизации задач с условной логикой, а также для одновременного управления несколькими большими языковыми моделями или инструментами.
Что такое граф в LangGraph?
Граф в общем случае — это нелинейная структура данных, состоящая из «вершин» (узлов) и «рёбер» (соединений между ними), которая отражает взаимосвязи между объектами.
В LangGraph именно эта структура графа используется для создания состоянийных, циклических рабочих процессов — таких, в которых ИИ может принимать решения, возвращаться к предыдущим шагам или разветвляться по разным путям в зависимости от промежуточных результатов.
LangChain против LangGraph
Архитектура проекта
(Конечная точка FastAPI + LangGraph + создание и сохранение транзакций)
Проблема
До внедрения LangGraph в финансовом приложении для создания транзакции требовалось напрямую вызывать конечную точку /transactions/add с заголовком данных, выглядевшим примерно так:
{
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"
}
Чтобы добраться до этой точки, потребовались два предварительных вызова API — один для получения списка categories и другой для получения payment_options — лишь бы получить идентификаторы, необходимые для передачи данных. Другими словами, создание одной транзакции представляло собой трехэтапный процесс, притом медленный.
Решение
Решением стало поручить ИИ-ассистенту выполнение всех трех этапов, в то время как пользователю достаточно описать простым языком, что он сделал со своими деньгами. С учетом этой цели вот как организована реализация.
Архитектура из трех этапов
- Конечная точка FastAPI
- Оркестратор LangGraph
- База данных PostgreSQL
1. Конечная точка FastAPI
Пользователь отправляет запрос на конечную точку FastAPI /assistance/transaction-entry с телом запроса, содержащим сообщение, описывающее операцию.
{
message: "Sent $5 to Rosy for Coffee through cash."
}
2. Orchestrator LangGraph
Оркестратор построен в виде графа, при этом каждый узел представляет операцию, а каждая рёбра — поток данных между операциями.
Первый узел, анализатор LLM, принимает сообщение пользователя и проверяет, содержит ли оно всю необходимую информацию для записи операции: является ли это доходом или расходом, какова сумма, в чём заключается цель операции, какой используется способ оплаты и так далее.
Если в сообщении уже содержатся все необходимые данные, LLM преобразует его в структурированные данные, например:
{
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
}
После удаления двух флаговых полей (is_complete и missing_info_message) эти структурированные данные передаются на узел Database Writer. Узел Database Writer вызывает метод create_transaction(), который записывает информацию о транзакции для пользователя в базе данных.
Но что произойдет, если в сообщении отсутствуют некоторые детали? Рассмотрим сообщение вроде этого:
{
message: "Sent $5 to Rosy for Coffee." // payment mode is not specified
}
Здесь способ оплаты не указан. В таком случае данные, извлеченные LLM, будут включать заполненные значения флагов, например:
{
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
}
Поскольку флаг is_complete здесь имеет значение False, сообщение missing_info_message направляется на другой узел, подключенный к анализатору LLM — узел Clarification. Этот путь активируется только тогда, когда значение is_complete равно False.
Узел Clarification получает missing_info_message и вызывает метод ask_again(), который возвращает этое сообщение в качестве ответа на первоначальный запрос FastAPI. Это означает окончание выполнения графа в текущей сессии — пользователь получает просто запрос на предоставление отсутствующей информации, в данном случае способа оплаты.
Предположим, что пользователь затем отправляет отсутствующую информацию, например:
{
message: "UPI"
}
Этот ответ приводит к повторной инициализации графа-оркестратора, который проходит ту же последовательность шагов, что и ранее.
Ключевое отличие этого второго прохода заключается в том, что никакая из ранее полученных информаций не теряется — история разговора сохраняется каждый раз, когда данные извлекаются большой языковой моделью (этот механизм сохранения рассматривается более подробно позже в статье). Поскольку теперь доступен параметр payment_option, который сочетается с ранее зафиксированными значениями, значение is_complete меняется на True, и готовые, отфильтрованные данные передаются в узел Database Writer в следующем виде:
{
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.
Затем узел Database Writer берет на себя обработку, имея в руках эти отфильтрованные данные. Напомним, что при ручном создании записи транзакции требовалось два дополнительных вызова API — один для categories и один для payment_options — лишь для получения соответствующих ID перед тем, как можно было создать саму транзакцию. Та же проблема возникает и здесь: в отфильтрованных данных содержатся фактические текстовые значения категорий и вариантов оплаты, а не их ID в базе данных, и база данных не примет сырые значения для этих полей.
Чтобы решить эту проблему, узел Database Writer должен выполнить запрос к базе данных для поиска соответствующих записей категорий и вариантов оплаты на основе значений, присутствующих в отфильтрованных данных.
Поскольку проект использует FastAPI вместе с SQLAlchemy, эти поиски реализуются в виде запросов SQLAlchemy.
Для 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
Кратко, логика выглядит так:
Выполняется запрос select для проверки наличия категории, указанной в data.category.
Если она уже существует, ID этой категории заменяет значение в data.category.
Если категории нет, создается новая запись категории, и вместо нее используется ее новый ID.
Тот же принцип применяется к 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
После того как определяются ID категории и варианта оплаты, объект данных полностью обновляется и готов к вставке; он выглядит следующим образом:
{
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"
}
С этими готовыми данными узел Database Writer вызывает метод create_transaction(), который фактически сохраняет запись транзакции в базе данных.
3. База данных PostgreSQL
Это представляет собой последнюю стадию архитектуры, на которой окончательные данные, переданные узлом Database Writer, записываются в таблицу transactions.
Результативная структура таблицы transactions выглядит следующим образом:
Реализация
Поскольку архитектура рассмотрена, пришло время подробно разобрать детали реализации этого оркестратора с использованием LangGraph. Обратите внимание, что здесь приведенный порядок не полностью соответствует порядку описания архитектуры выше. Вместо этого реализация организована следующим образом:
- LangGraph Orchestrator
- База данных PostgreSQL
- Конечная точка FastAPI
1. LangGraph Orchestrator
Сам оркестратор находится в файле src/assistance/graph.py. Этот файл отвечает за настройку LLM, определение узлов графа, создание связей между этими узлами и, наконец, сборку всего в работоспособный граф.
Как уже упоминалось, этот оркестратор состоит из трех узлов: Analyzer LLM, Database Writer и узел Clarification.
Узел Analyzer LLM (Groq)
Этот узел по сути представляет собой языковую модель, задача которой — понять намерение пользователя и проверить, содержит ли сообщение все необходимые и корректные детали. Вместо того чтобы создавать собственную модель с нуля, в этом проекте используется Groq для выполнения сложных задач.
Что такое Groq?
Groq — это фреймворк на Python с открытым исходным кодом, предназначенный для работы с данными в графовой структуре. Он предоставляет разработчикам мощные инструменты для запросов, фильтрации и агрегации информации, хранящейся в виде графов, и отлично подходит для обработки крупных наборов графовых данных, таких как социальные сети, знаниевые графы или системы рекомендаций, как описано в статье GeekForGeeks об API Groq.
С помощью хостингового API Groq вы можете отправлять запросы к широко используемым открытым моделям — openai/gpt-oss-120b используется в этом проекте — и получать ответы, которые приходят значительно быстрее, чем обычно от других поставщиков аналогичных моделей.
Почему Groq?
Грок был выбран вместо таких альтернатив, как ChatOpenAI или ChatAnthropic, по нескольким причинам:
- Скорость: Грок использует специально разработанное оборудование под названием LPU (ядра обработки языка) вместо GPU, от которых зависят большинство других поставщиков, что позволяет достигать высокой скорости обработки запросов.
- Полезный бесплатный тариф: бесплатный тариф Грока достаточно щедр, чтобы поддерживать индивидуальные или образовательные проекты, не создавая значительных затрат на API во время экспериментов.
- Совместимость через LangChain: класс
langchain_groq.ChatGroqинтегрируется с LangChain и LangGraph так же, как это делаютChatOpenAIилиChatAnthropic. Это означает, что в будущем при смене поставщика не потребуется перерабатывать логику работы графа — достаточно будет заменить клиентское приложение.
Как получить ключ API Groq
Groq позволяет генерировать бесплатные ключи API для разработки. Вот как их получить:
- Перейдите на https://console.groq.com и войдите или зарегистрируйтесь.
- В навигационной панели выберите опцию API Keys.
- Выберите вариант Создать ключ API.
- Появится форма с запросом указать имя (для этого проекта использовано
transaction-assistant) и срок действия ключа. После заполнения формы нажмите Отправить. - Ключ отображается только один раз, сразу после создания — поэтому обязательно скопируйте его сразу же.
После создания все ваши ключи появятся в основном списке на этой странице.
Использование ключа API Groq в коде FastAPI
Добавьте ключ API Groq в файл .env, расположенный в корне проекта, рядом с другими переменными окружения:
GROQ_API_KEY = "gsk_***************************************DyxM"
Существует несколько способов загрузки переменных окружения в те модули, которые ими нуждаются. В этом проекте используется специальный класс настроек:
Определите класс Settings в файле 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()
Затем импортируйте этот объект настроек там, где он требуется:
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.
Настройка LLM
Прежде чем настроить LLM, установите LangGraph и LangChain вместе с интеграцией Groq:
pip install -U langgraph langchain langchain-groq
Затем создается экземпляр клиента Groq и настраивается с использованием конкретной модели:
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)
Здесь ChatGroq выступает в роли обертки LangChain вокруг чат-моделей Groq, позволяя взаимодействовать с ними через стандартный интерфейс LangChain вместо ручного формирования HTTP-запросов.
assistance_llm = ChatGroq(model="openai/gpt-oss-120b", temperature=0.2,
api_key=settings.GROQ_API_KEY)
В этом фрагменте создается упомянутая выше инстанция клиента Groq, настроенная с выбранным моделью и значением температуры, а также авторизованная с использованием API-ключа, взятого из настроек окружения.
Температура — это параметр, обычно варьирующийся от 0 до 1, который определяет степень случайности или творческости ответов модели. Более высокое значение, например 0.8, способствует получению более разнообразных и креативных результатов, в то время как более низкое значение, например 0.2, делает ответы более строгими и предсказуемыми. В этом проекте установлено temperature = 0.2.
structured_llm = assistance_llm.with_structured_output(
ExtractedTransactionSchema)
Этот код оборачивает LLM таким образом, что вместо возврата простого текста он генерирует объект на Python, полностью соответствующий ExtractedTransactionSchema. Внутренне это достигается путем указания модели на генерацию вывода, соответствующего данной схеме, а затем автоматической обработки и проверки этого вывода — что устраняет необходимость вручную интерпретировать сырой текст модели.
Сама структура ExtractedTransactionSchema определена в файле 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
Обратите внимание, что на этом этапе LLM ещё не вызывался — эта операция лишь определяет формат вывода, который должен быть получен после его вызова.
Состояние графа
Состояние графа представляет собой структуру данных, которая циркулирует внутри графа и обновляется при его работе. Можно представить его как рабочую память дирижера: она хранит всю информацию, которую отслеживает граф и изменяет по мере выполнения каждого шага. Для этого помощника транзакций состояние графа определяется следующим образом:
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
Давайте разберем, что на самом деле делает этот код:
from pydantic import BaseModel, Field
Pydantic — это библиотека для проверки данных, используемая здесь. BaseModel — это родительский класс, от которого наследуется при определении структурированного объекта с проверкой типов, такого как GraphState. Field позволяет привязывать метаданные — описания, значения по умолчанию и т. д. — к каждому отдельному атрибуту.
from typing import Annotated, List, Optional
import operator
Это инструменты типизации в Python. Optional указывает на то, что поле может быть пустым и содержать значение None. List означает, что атрибут представляет собой список элементов. Annotated, используемый вместе с operator.add, сообщает LangGraph: «когда узел возвращает новое значение для этого поля, добавьте его к уже существующему содержимому вместо замены». Именно этот механизм позволяет conversation_history расти с каждым новым сообщением, вместо того чтобы очищаться после каждой новой передачи данных.
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: самое последнее сообщение, отправленное пользователем в ходе текущей операции.conversation_history: полный список всех предыдущих сообщений, накапливаемый по мере взаимодействия без их перезаписи.
extracted: заполняется после того, как большая языковая модель извлекла структурированные данные о транзакции из разговора. Изначально его значение — None, поскольку на начальном этапе выполнения ничего ещё не было извлечено.final_response: сообщение, которое в итоге отправляется пользователю — либо подтверждение о том, что транзакция зарегистрирована, либо дополнительный вопрос с просьбой предоставить больше деталей.Запрос на извлечение
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"]
)
Это буквальная инструкция, передаваемая в LLM на естественном языке — она указывает, какие поля необходимо искать, что делать при отсутствии информации и как должна быть структурирована ответная информация. Поскольку structured_llm уже обеспечивает соблюдение схемы на уровне вывода, задача промпта заключается в основном в направлении мышления модели: определении того, что означает «полная информация», формулировке уточняющих вопросов и т. д., в то время как схема отвечает за форматирование.
Экстрактор
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]}
Функция extractor выполняет следующее:
- Сливает все предыдущие сообщения с текущим, чтобы LLM мог видеть полный контекст.
- Передает полученный объединённый текст в LLM.
- Возвращает структурированный объект
ExtractedTransactionSchema.
operator.add, настроенному для этого поля.Решение
route_after_extraction
def route_after_extraction(state: GraphState):
return "create_transaction" if state.extracted.is_complete else "ask_again"
Эта функция не выполняет никакой реальной обработки — её единственная задача — принять решение. В зависимости от того, пометил ли LLM извлеченные данные как полные, она возвращает строку, указывающую графу, какой узел следует обработать дальше. Можно считать её логикой разветвления в диаграмме потока: граф анализирует значение, возвращаемое этой функцией, и следует соответствующему пути — либо к функции create_transaction для записи транзакции, либо к функции ask_again для запроса дополнительной информации.
Узел записи в базу данных
create_transaction_node
Этот узел отвечает за запись завершённых данных транзакции в базу данных от имени соответствующего пользователя. Функция create_transaction_node, реализующая узел DB Writer, выглядит следующим образом:
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}
Это слишком много информации за один раз, поэтому давайте рассмотрим её по частям.
from langchain_core.runnables import RunnableConfig
Тип, заменяющий объект config, передаваемый в любой узел. Он существует исключительно в качестве указателя типа, поэтому любой, кто читает сигнатуру create_transaction_node, сразу понимает, какой формы имеет config.
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.ext.asyncio import AsyncSession
Это стандартные импорты SQLAlchemy, необходимые для обработки ошибок базы данных и для указания типа асинхронной сессии базы данных, используемой для взаимодействия с PostgreSQL.
from src.transaction import controller
from src.transaction.schema import TransactionCreateSchema
Здесь используется логика создания транзакций, уже применявшаяся в других частях приложения, вместе с соответствующей схемой входных данных. Использование этой логики позволяет ассистенту создавать транзакции по точно такому же кодовому пути, как и обычный CRUD API, вместо того чтобы дублировать эту логику здесь.
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
Это вспомогательные элементы, предназначенные для преобразования названий категорий и вариантов оплаты, выделенных большой языковой моделью, в реальные строки и ID в базе данных, а также для создания новых записей, если они ещё не существуют.
Примечание: логика поиска/создания для categories и payment_options была объединена в один универсальный помощник get_or_create, поскольку оба модели требовали практически одинакового поведения.
async def create_transaction_node(state: GraphState, config: RunnableConfig):
Функция create_transaction_node выполняется только после подтверждения полноты извлеченных данных. Она объявлена как async, поскольку выполняет реальные операции с базой данных, и принимает параметры config и state, чтобы иметь доступ к активной сессии базы данных и залогиненному пользователю. Эти два значения поступают из API-маршрута, а не от LLM или состояния разговора, поскольку они относятся к конкретному запросу, а не к текущему диалогу.
session: AsyncSession = config["configurable"]["session"]
user = config["configurable"]["user"]
data = state.extracted
Этот код извлекает сессию, пользователя и данные транзакции, полученные в результате обработки.
try:
category_id = await get_or_create(...)
payment_option_id = await get_or_create(...)
Поскольку большая языковая модель извлекла только названия категории и способа оплаты, такие как «Продукты питания» или «UPI», а не их идентификаторы в базе данных, на этом этапе проверяется наличие соответствующей записи для текущего пользователя. Если такой записи нет, она создается. В любом случае возвращается соответствующий идентификатор.
payload = TransactionCreateSchema(...)
await controller.create_transaction(payload, session, user)
await session.commit()
Полезная нагрузка формируется в том же формате, который ожидается существующей логикой создания транзакций, затем передается в ту же функцию-контроллер, что позволяет повторно использовать существующую логику приложения вместо её переписывания. После этого транзакция в базе данных фиксируется для сохранения изменений.
except SQLAlchemyError as err:
await session.rollback()
print(...)
return {"final_response": "Something went wrong..."}
Если происходит сбой на уровне базы данных, все частичные изменения откатываются, и возвращается понятное сообщение об ошибке вместо того, чтобы запрос завершился сбоем. Это предотвращает ситуацию, когда, например, создается новая категория, но при этом отсутствует соответствующая транзакция.
message = f"Added {data.transaction_type} of {data.amount} under '{data.category}' ({data.payment_option})"
return {"final_response": message}
При успешном выполнении формируется сообщение подтверждения, понятное человеку, которое затем возвращается для обновления состояния.
Узел уточнения
ask_again
def ask_again_node(state: GraphState):
return {"final_response": state.extracted.missing_info_message}
Это простой запасной вариант. Как описано ранее, этот узел запускается только тогда, когда параметр is_complete у извлеченных данных равен False, при этом в поле missing_info_message указывается полезное сообщение.
Внутри ask_again_node функция получает параметр state, что позволяет ей обращаться к свойствам state.extracted.is_complete и state.extracted.missing_info_message.
Короче говоря, когда отсутствует какая-либо информация, этот узел просто пересылает уточняющий вопрос, уже сформулированный большой языковой моделью во время извлечения данных, чтобы пользователь точно знал, что нужно предоставить дальше.
Сборка графа
from langgraph.graph import StateGraph, START, END
graph_builder = StateGraph(GraphState)
Здесь создается новый конструктор графа, которому указывается, что каждый узел в графе будет читать и записывать данные в объект, соответствующий структуре 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")
Здесь задается точка входа: каждая экспекуция графа начинается с узла extractor.
graph_builder.add_conditional_edges(
"extractor",
route_after_extraction,
{
"create_transaction": "create_transaction",
"ask_again": "ask_again",
},
)
Здесь происходит разветвление. Как только extractor завершает работу, LangGraph вызывает route_after_extraction, чтобы определить следующий шаг. Независимо от того, какую строку он возвращает — create_transaction или ask_again — эта строка ищется в этой таблице сопоставлений, которая связывает каждую строку решения с конкретным узлом, к которому следует перейти.
graph_builder.add_edge("create_transaction", END)
graph_builder.add_edge("ask_again", END)
Обе возможные ветви прерывают выполнение графа по завершении своей работы, когда достигается точка END по любому из путей.
Компиляция с учетом использования памяти
from langgraph.checkpoint.memory import MemorySaver
memory = MemorySaver()
assistance_graph = graph_builder.compile(checkpointer=memory)
Вызов функции compile() преобразует определение графа в объект, способный к выполнению. Параметр checkpointer=memory включает механизм сохранения состояния, описанный ранее, что позволяет при повторном вызове графа с тем же thread_id возобновить работу с того момента, где она была прервана, вместо её полного перезапуска.
Итоговый код
Теперь слой оркестрации готов. Вот полностью завершённый файл (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. База данных PostgreSQL
Взаимодействие с базой данных на этом этапе уже реализовано внутри функции create_transaction_node, рассмотренной ранее, где окончательные данные транзакции записываются в таблицу transactions.
3. Конечная точка 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"]}
Клиенты обращаются к этому эндпоинту (/assistance/transaction-entry) и в теле запроса указывают сообщение с описанием сути операции.
Давайте рассмотрим функцию каждого элемента.
@assistance_routes.post("/transaction-entry", status_code=status.HTTP_201_CREATED)
Здесь настраивается маршрут POST по адресу /transaction-entry. Установка параметра status_code=status.HTTP_201_CREATED указывает FastAPI, какой код статуса возвращать по умолчанию при успешном выполнении. 201 — это стандартный код, означающий «был создан новый ресурс», что подходит здесь, поскольку успешный запрос приводит к созданию новой строки транзакции.
Сборка конфигурации графа:
config = {
"configurable": {
"thread_id": str(user.id),
"session": session,
"user": user
}
}
Здесь формируется объект config, который передается при вызове графа. Объект config содержит значения, специфичные для запроса, которые не должны храниться в самом состоянии диалога.
"thread_id": str(user.id): этое значение используется чекпойнтером LangGraph для определения истории разговора того пользователя, данные которого необходимо получить и обновить. Благодаря привязке к ID аутентифицированного пользователя у каждого пользователя автоматически формируется изолированный, постоянный поток, поэтому незавершённая запись транзакции одного пользователя никогда не может повлиять на другого. Это значение преобразуется в строку, поскольку чекпойнтер ожидаетthread_idв виде строки, тогда какuser.idобычно представляет собой UUID."session"и"user": эти параметры передаются дальше, чтобы функцияcreate_transaction_node, выполняющаяся внутри графа, имела доступ к активной сессии базы данных и к информации о том, кто отправляет запрос.
Вызов графа:
result = await assistance_graph.ainvoke(
{"user_input": payload.message}, config=config)
Это строка, которая фактически запускает выполнение. ainvoke является асинхронным аналогом запуска графа; использование синхронного invoke приведёт к блокировке цикла событий, что важно здесь, поскольку create_transaction_node внутренне выполняет асинхронные операции с базой данных.
- Первый аргумент,
{"user_input": payload.message}, обозначает начальное состояние для данной работы. Требуется явно указать толькоuser_input; остальные поляGraphState(conversation_history,extracted,final_response) имеют значения по умолчанию или заполняются по мере выполнения алгоритма. Если у существующегоthread_idуже сохраняется история, LangGraph объединяет этот новый ввод с уже хранящимся состоянием, вместо того чтобы начинать с нуля. config=configпередаёт всё, что было подготовлено на предыдущем этапе:thread_idдля поиска нужного состояния, а такжеsessionиuser— узлы, отвечающие за запись в базу данных.
await приостанавливает работу этого корутинного функционала до завершения полной обработки графа, поскольку ainvoke возвращает корутину, которую необходимо ожидать перед тем, как результат станет использоваемым.То, что возвращается как result, — это окончательное состояние GraphState, представленное в виде словаря, отражающего ход выполнения графа независимо от того, прервалось ли оно на этапе create_transaction или на этапе ask_again.
Возврат ответа:
return {"response": result["final_response"]}
Маршрут завершается возвратом простого словаря, содержащего только текст окончательного ответа. FastAPI сам преобразует его в JSON-пакет для клиента, получая результат примерно такого вида:
{ "response": "Added expense of 450 under 'Groceries' (UPI)" }
Этот текст идентичен тому, что было сгенерировано ранее в функциях create_transaction_node или ask_again_node. Сам маршрут не знает, какая именно ветка была выполнена; он просто передает содержимое переменной final_response.
Заключение
Работа с этим помощником по транзакциям показывает то, что обычно упускается в учебных материалах: сложность реализации функций ИИ заключается не в том, чтобы попросить модель дать ответ, а в обеспечении безопасного поведения этого ответа при взаимодействии с реальной системой. Написание запросов к модели — это простая часть. Настоящие инженерные усилия требуются для создания схем, обеспечивающих структурированный вывод, условных алгоритмов, определяющих, следует ли сохранять данные или запросить уточнения, а также для правильного хранения состояния на протяжении нескольких этапов диалога.
LangGraph оказался подходящим инструментом именно для этого проекта, поскольку рабочий процесс требовал реального принятия решений, а не простого прохождения данных от входа к выходу. Если ваше решение требует линейного потока, то простой вызов LLM или цепочка LangChain, скорее всего, будут более простыми и подходящими инструментами. Однако когда логика ИИ нуждается в разветвлениях, сохранении информации или паузе для сбора дополнительных данных перед продолжением работы, графовая структура перестаёт казаться избыточной сложностью и становится наиболее разумным способом моделирования такого потока.
Ссылки
Похожие материалы
- LangChain против LangGraph: выбор между цепочками и графами с состоянием — Узнайте, в чём различия между линейными строительными блоками LangChain и работающими с состоянием ветвящимися рабочими процессами LangGraph, а также как определить, что лучше подойдёт для вашего ИИ-приложения.
- Создание ИИ-агента с нуля: шаблоны, ReAct и LangGraph — Ознакомьтесь с основными концепциями ИИ-агентов — планированием, использованием инструментов, рефлексией и шаблоном ReAct — а также с тем, как LangChain и LangGraph помогают вручную создать такой агент.