RAG на Elasticsearch и Cohere

Урок 1 из 114 курса «Cohere: LLM University»: официальный курс Cohere LLM University (Кохир) на русском языке. Этот урок бесплатный.

Комплексный RAG с использованием Elasticsearch и Cohere

Узнайте, как использовать Inference API для семантического поиска и использовать API Cohere для RAG.

🧰 Требования

Для этого примера вам понадобится:

Примечание: Хотя этот учебник интегрирует Cohere с бессерверным проектом Elastic Cloud, вы также можете интегрировать его с вашим самостоятельно управляемым развертыванием Elasticsearch или развертыванием Elastic Cloud, просто переключившись с использования Serverless endpoint в клиенте Elasticsearch.

Создайте развертывание Elastic Serverless

Если у вас нет развертывания Elastic Cloud, зарегистрируйтесь здесь для получения бесплатной пробной версии и запросите доступ к Elastic Serverless

Установите пакеты и подключитесь к клиенту Elasticsearch Serverless

Для начала нам нужно подключиться к нашему развертыванию Elastic Serverless с помощью клиента Python.

Сначала нам нужно установить следующие пакеты с помощью pip:

  • elasticsearch_serverless
  • cohere

После установки на панели управления 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.

Полезные гиды