Комплексный RAG с использованием Elasticsearch и Cohere
Узнайте, как использовать Inference API для семантического поиска и использовать API Cohere для RAG.
🧰 Требования
Для этого примера вам понадобится:
Аккаунт Elastic Serverless через Elastic Cloud, доступный с бесплатной пробной версией
Аккаунт Cohere с производственным API key
Python 3.7 или выше
Примечание: Хотя этот учебник интегрирует Cohere с бессерверным проектом Elastic Cloud, вы также можете интегрировать его с вашим самостоятельно управляемым развертыванием Elasticsearch или развертыванием Elastic Cloud, просто переключившись с использования Serverless endpoint в клиенте Elasticsearch.
Создайте развертывание Elastic Serverless
Если у вас нет развертывания Elastic Cloud, зарегистрируйтесь здесь для получения бесплатной пробной версии и запросите доступ к Elastic Serverless
Установите пакеты и подключитесь к клиенту Elasticsearch Serverless
Для начала нам нужно подключиться к нашему развертыванию Elastic Serverless с помощью клиента Python.
Сначала нам нужно установить следующие пакеты с помощью pip:
elasticsearch_serverlesscohere
После установки на панели управления Serverless найдите свой endpoint URL и создайте свой API key.
pip install elasticsearch_serverless cohere
Далее нам нужно импортировать необходимые модули. 🔐 ПРИМЕЧАНИЕ: getpass позволяет нам безопасно запрашивать у пользователя учетные данные, не выводя их в терминал и не сохраняя в памяти.
from elasticsearch_serverless import Elasticsearch, helpers
from getpass import getpass
import cohere
import json
import requests
Теперь мы можем создать экземпляр клиента Python Elasticsearch.
Сначала мы запрашиваем у пользователя его endpoint и закодированный API key.
Затем мы создаем объект client, который является экземпляром класса Elasticsearch.
При создании вашего Elastic Serverless API key убедитесь, что вы включили Control security privileges и отредактировали cluster privileges, чтобы указать "cluster": ["all"]
ELASTICSEARCH_ENDPOINT = getpass("Elastic Endpoint: ")
ELASTIC_API_KEY = getpass("Elastic encoded API key: ") # Используйте закодированный API key
client = Elasticsearch(
ELASTICSEARCH_ENDPOINT,
api_key=ELASTIC_API_KEY
)
Подтвердите, что клиент подключен, с помощью этого теста:
print(client.info())
Создайте inference endpoint
Давайте создадим inference endpoint, используя Create inference API.
Для этого вам понадобится Cohere API key, который вы можете найти в своем аккаунте Cohere в разделе API keys. Для выполнения шагов в этом ноутбуке требуется production key, так как использование бесплатной пробной версии Cohere API ограничено.
COHERE_API_KEY = getpass("Enter Cohere API key: ")
# Удалить модель вывода, если она уже существует
client.options(ignore_status=[404]).inference.delete_model(inference_id="cohere_embeddings")
client.inference.put_model(
task_type="text_embedding",
inference_id="cohere_embeddings",
body={
"service": "cohere",
"service_settings": {
"api_key": COHERE_API_KEY,
"model_id": "embed-v4.0",
"embedding_type": "int8",
"similarity": "cosine"
},
"task_settings": {},
},
)
Создайте ingest pipeline с inference processor
Создайте ingest pipeline с inference processor, используя метод put_pipeline. Ссылайтесь на inference endpoint, созданный выше, как на model_id для вывода данных, которые поступают в pipeline.
# Удалить конвейер приема, если он уже существует
client.options(ignore_status=[404]).ingest.delete_pipeline(id="cohere_embeddings")
client.ingest.put_pipeline(
id="cohere_embeddings",
description="Ingest pipeline for Cohere inference.",
processors=[
{
"inference": {
"model_id": "cohere_embeddings",
"input_output": {
"input_field": "text",
"output_field": "text_embedding",
},
}
}
],
)
Давайте отметим несколько важных параметров из этого вызова API:
inference: Процессор, который выполняет вывод с использованием модели машинного обучения.model_id: Указывает ID inference endpoint, который будет использоваться. В этом примере ID модели установлен наcohere_embeddings.input_output: Указывает входные и выходные поля.input_field: Имя поля, из которого создается представлениеdense_vector.output_field: Имя поля, которое содержит результаты вывода.
Создайте индекс
Должно быть создано сопоставление destination index – индекса, который содержит embeddings, которые модель создаст на основе вашего входного текста. Destination index должен иметь поле с типом поля dense_vector для индексации вывода модели Cohere.
Давайте создадим индекс с именем cohere-wiki-embeddings с необходимыми нам сопоставлениями.
client.indices.delete(index="cohere-wiki-embeddings", ignore_unavailable=True)
client.indices.create(
index="cohere-wiki-embeddings",
settings={"index": {"default_pipeline": "cohere_embeddings"}},
mappings={
"properties": {
"text_embedding": {
"type": "dense_vector",
"dims": 1024,
"element_type": "byte"
},
"text": {"type": "text"},
"wiki_id": {"type": "integer"},
"url": {"type": "text"},
"views": {"type": "float"},
"langs": {"type": "integer"},
"title": {"type": "text"},
"paragraph_id": {"type": "integer"},
"id": {"type": "integer"}
}
},
)
Вставьте документы
Давайте вставим наш пример wiki dataset. Для выполнения этого шага вам потребуется production аккаунт Cohere, иначе прием документации будет прерван по таймауту из-за ограничений скорости запросов API.
url = "https://raw.githubusercontent.com/cohere-ai/cohere-developer-experience/main/notebooks/data/embed_jobs_sample_data.jsonl"
response = requests.get(url)
# Загрузить данные ответа в объект JSON
jsonl_data = response.content.decode('utf-8').splitlines()
# Подготовить документы для индексации
documents = []
for line in jsonl_data:
data_dict = json.loads(line)
documents.append({
"_index": "cohere-wiki-embeddings",
"_source": data_dict,
}
)
# Использовать bulk endpoint для индексации
helpers.bulk(client, documents)
print("Done indexing documents into `cohere-wiki-embeddings` index!")
Гибридный поиск
После того как dataset был обогащен embeddings, вы можете запрашивать данные с помощью hybrid search.
Передайте query_vector_builder в k-nearest neighbor (kNN) vector search API и предоставьте query text и модель, которую вы использовали для создания embeddings.
query = "When were the semi-finals of the 2022 FIFA world cup played?"
response = client.search(
index="cohere-wiki-embeddings",
size=100,
knn={
"field": "text_embedding",
"query_vector_builder": {
"text_embedding": {
"model_id": "cohere_embeddings",
"model_text": query,
}
},
"k": 10,
"num_candidates": 50,
},
query={
"multi_match": {
"query": query,
"fields": ["text", "title"]
}
}
)
raw_documents = response["hits"]["hits"]
# Отобразить первые 10 результатов
for document in raw_documents[0:10]:
print(f'Title: {document["_source"]["title"]}\nText: {document["_source"]["text"]}\n')
# Форматировать документы для ранжирования
documents = []
for hit in response["hits"]["hits"]:
documents.append(hit["_source"]["text"])
Ранжирование
Чтобы эффективно объединить результаты нашего vector и BM25 retrieval, мы можем использовать модель Rerank 3 от Cohere через inference API для обеспечения окончательного, более точного, semantic reranking наших результатов.
Сначала создайте inference endpoint с вашим Cohere API key. Убедитесь, что вы указали имя для вашего endpoint и model_id одной из моделей rerank. В этом примере мы будем использовать Rerank 3.
# Удалить модель вывода, если она уже существует
client.options(ignore_status=[404]).inference.delete_model(inference_id="cohere_rerank")
client.inference.put_model(
task_type="rerank",
inference_id="cohere_rerank",
body={
"service": "cohere",
"service_settings":{
"api_key": COHERE_API_KEY,
"model_id": "rerank-english-v3.0"
},
"task_settings": {
"top_n": 10,
},
}
)
Теперь вы можете переранжировать свои результаты, используя этот inference endpoint. Здесь мы передадим query, который мы использовали для retrieval, вместе с документами, которые мы только что получили с помощью hybrid search.
Inference service ответит списком документов в порядке убывания релевантности. Каждый документ имеет соответствующий index (отражающий порядок документов при отправке на inference endpoint), и если параметр задачи “return_documents” установлен в True, то тексты документов также будут включены.
В этом случае мы установим response в False и реконструируем входные документы на основе index, возвращенного в response.
response = client.inference.inference(
inference_id="cohere_rerank",
body={
"query": query,
"input": documents,
"task_settings": {
"return_documents": False
}
}
)
# Реконструировать входные документы на основе индекса, предоставленного в ответе rerank
ranked_documents = []
for document in response.body["rerank"]:
ranked_documents.append({
"title": raw_documents[int(document["index"])]["_source"]["title"],
"text": raw_documents[int(document["index"])]["_source"]["text"]
})
# Вывести первые 10 результатов
for document in ranked_documents[0:10]:
print(f"Title: {document['title']}\nText: {document['text']}\n")
Retrieval augmented generation
Теперь, когда мы ранжировали наши результаты, мы можем легко превратить это в RAG систему с помощью Cohere's Chat API. Передайте полученные документы вместе с query и посмотрите grounded response, используя новейшую генеративную модель Cohere Command R+.
Сначала мы создадим Cohere client.
co = cohere.Client(COHERE_API_KEY)
Далее мы можем легко получить grounded generation с citations из Cohere Chat API. Мы просто передаем user query и документы, полученные из Elastic, в API и выводим наш grounded response.
response = co.chat(
message=query,
documents=ranked_documents,
model='command-r-plus'
)
source_documents = []
for citation in response.citations:
for document_id in citation.document_ids:
if document_id not in source_documents:
source_documents.append(document_id)
print(f"Query: {query}")
print(f"Response: {response.text}")
print("Sources:")
for document in response.documents:
if document['id'] in source_documents:
print(f"{document['title']}: {document['text']}")
И вот оно! Быстрая и простая реализация hybrid search и RAG с Cohere и Elastic.