在今天的文章中我们以一个 IMDB 电影数据为例来展示如何把数据写入到 Elasticsearch 中,并对它进行 AI 分析。

更多阅读:如何通过 Claude Code 来写入 CSV 数据到 Elasticsearch

准备工作

安装 Elasticsearch 及 Kibana

如果你还没有安装好自己的 Elasticsearch 及 Kibana,你可以启动云服务,比如阿里云,腾讯云,或者 Elastic Cloud。你可以参观看如下的视频:

建议你观看上面的 “为 AI 安装自己的 Elasticsearch 集群” 来进行安装。

为了能够使我们使用向量搜索对 overview 字段进行搜索,我们使用 E5 多语言模型。我们还需要进行如下的安装:

我们需要等后一点时间来等待下载完成。在安装过程中,如果由于任何原因导致安装不成功,你也可以直接使用如下的命令来删除它,然后再次下载:

DELETE /_ml/trained_models/.multilingual-e5-small-elasticsearch?force=true

一旦下载完毕后,我们可以对它进行测试:

GET _inference/_all

我们可以看到一个叫做 .multilingual-e5-small-elasticsearch 的 inference_id。我们可以做如下的一个简单的测试:

POST _inference/.multilingual-e5-small-elasticsearch
{
  "input": "The sky above the port was the color of television tuned to a dead channel."
}

很显然,这个推理端点工作正常。它可以帮我们把文字转换为向量。

下载 IMDB 数据

这个数据集可以在地址进行下载:

等下载完毕后,我们可以把它放到我们的目录中:

$ pwd
/Users/liuxg/python/imdb
$ ls
imdb_movies.csv.zip

我们接下来把它进行解压缩:

$ unzip imdb_movies.csv.zip 
Archive:  imdb_movies.csv.zip
  inflating: imdb_movies.csv         
$ ls
imdb_movies.csv     imdb_movies.csv.zip

我们可以看到上面的被解压缩的文件 imdb_movies.csv。文件的结构如下:

很显然,我们可以看到一个叫做 names, date_x, score, genre, overview 等字段。

写入数据到 Elasticsearch

我们有几种方法来写入数据。我们也可以仿照文章 “Elasticsearch:Retrievers 介绍 - Python Jupyter notebook” 来写入我们的数据。

使用 Kibana 上传写入

我们打开 Kibana:

经过调整后的 mapping 是这样的:

{
  "properties": {
    "budget_x": {
      "type": "double"
    },
    "country": {
      "type": "keyword"
    },
    "crew": {
      "type": "text"
    },
    "date_x": {
      "type": "date",
      "format": "MM/dd/yyyy||MM/dd/yyyy[ ]"
    },
    "genre": {
      "type": "keyword"
    },
    "names": {
      "type": "text"
    },
    "orig_lang": {
      "type": "keyword"
    },
    "orig_title": {
      "type": "text"
    },
    "overview": {
      "type": "text",
      "copy_to": "overview_semantic"
    },
    "overview_semantic": {
      "type": "semantic_text",
      "inference_id": ".multilingual-e5-small-elasticsearch"
    },
    "revenue": {
      "type": "double"
    },
    "score": {
      "type": "double"
    },
    "status": {
      "type": "keyword"
    }
  }
}

完美。我们已经成功地写入了所有的数据。我们可以直接点击 Agent Builder 对数据进行查询。

当然,在进行之前,我们也可以进入到 Kibana 查看数据:

GET imdb/_mapping

查看我们的 mapping:

{
  "imdb": {
    "mappings": {
      "_meta": {
        "created_by": "file-data-visualizer"
      },
      "properties": {
        "budget_x": {
          "type": "double"
        },
        "country": {
          "type": "keyword"
        },
        "crew": {
          "type": "text"
        },
        "date_x": {
          "type": "keyword"
        },
        "genre": {
          "type": "keyword"
        },
        "names": {
          "type": "text"
        },
        "orig_lang": {
          "type": "keyword"
        },
        "orig_title": {
          "type": "text"
        },
        "overview": {
          "type": "text",
          "copy_to": [
            "overview_semantic"
          ]
        },
        "overview_semantic": {
          "type": "semantic_text",
          "inference_id": ".multilingual-e5-small-elasticsearch",
          "model_settings": {
            "service": "elasticsearch",
            "task_type": "text_embedding",
            "dimensions": 384,
            "similarity": "cosine",
            "element_type": "float"
          }
        },
        "revenue": {
          "type": "double"
        },
        "score": {
          "type": "double"
        },
        "status": {
          "type": "keyword"
        }
      }
    }
  }
}

也可以针对数据进行查询:

这样我们的数据就已经写入完成了。

使用 Python 代码写入数据

针对有些开发者不是很喜欢使用 Upload file 上面的方法,你可以使用 Python 代码来实现上面的写入。针对我们的数据,我们可以有如下的代码来实现。在运行下面的代码之前,我们先运行下面的命令来删除已经写入的数据:

DELETE imdb

我们先在当前目录创建一个叫做 .env 的文件:

$ pwd
/Users/liuxg/python/imdb
$ ls -al
total 19328
drwxr-xr-x   6 liuxg  staff      192 Jul  9 11:38 .
drwxr-xr-x@ 62 liuxg  staff     1984 Jul  9 09:48 ..
-rw-r--r--   1 liuxg  staff      105 Jul  9 11:34 .env
-rw-r--r--@  1 liuxg  staff  6723214 Apr 28  2023 imdb_movies.csv
-rw-r--r--@  1 liuxg  staff  2980385 Jul  9 09:51 imdb_movies.csv.zip
-rw-r--r--   1 liuxg  staff     4356 Jul  9 11:38 ingest_imdb.py
ES_URL="https://localhost:9200"
ES_API_KEY="ckFQc1JKOEJHN2xFUjlaTlNKbUw6alU5MDgxOUlmM0hjLTE0WXdyQVNJUQ=="

你需要根据自己的配置修改上面的 ES_URL 地址来指向自己的 Elasticsearch 端点。如果大家还不是很清楚如何得到 ES_API_KEY 的话,建议大家观看我之前的视频 “ Elastic 黑客松比赛说明”。我简单展示如下:

拷贝上面的 API key 即可:ckFQc1JKOEJHN2xFUjlaTlNKbUw6alU5MDgxOUlmM0hjLTE0WXdyQVNJUQ==

这样,我们的 .env 配置就完成了。

ingest_imdb.py

#!/usr/bin/env python3
"""Ingest imdb_movies.csv into Elasticsearch, using connection settings from .env."""

import csv
import os
import sys
import urllib3

from dotenv import load_dotenv
from elasticsearch import Elasticsearch
from elasticsearch.helpers import bulk, BulkIndexError
from elastic_transport import TlsError

INDEX_NAME = "imdb"
CSV_PATH = os.path.join(os.path.dirname(os.path.abspath(__file__)), "imdb_movies.csv")
BULK_CHUNK_SIZE = 50
REQUEST_TIMEOUT = 300

INDEX_MAPPING = {
    "mappings": {
        "properties": {
            "budget_x": {"type": "double"},
            "country": {"type": "keyword"},
            "crew": {"type": "text"},
            "date_x": {"type": "keyword"},
            "genre": {"type": "keyword"},
            "names": {"type": "text"},
            "orig_lang": {"type": "keyword"},
            "orig_title": {"type": "text"},
            "overview": {"type": "text", "copy_to": ["overview_semantic"]},
            "overview_semantic": {
                "type": "semantic_text",
                "inference_id": ".multilingual-e5-small-elasticsearch",
                "model_settings": {
                    "service": "elasticsearch",
                    "task_type": "text_embedding",
                    "dimensions": 384,
                    "similarity": "cosine",
                    "element_type": "float",
                },
            },
            "revenue": {"type": "double"},
            "score": {"type": "double"},
            "status": {"type": "keyword"},
        },
    }
}


def build_client(es_url: str, es_api_key: str) -> Elasticsearch:
    """Connect to Elasticsearch, working for both trusted and self-signed TLS certs."""
    try:
        client = Elasticsearch(
            es_url, api_key=es_api_key, verify_certs=True, request_timeout=REQUEST_TIMEOUT
        )
        client.info()
        return client
    except TlsError:
        print("Certificate could not be verified (self-signed?), retrying with verify_certs=False", file=sys.stderr)
        urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
        client = Elasticsearch(
            es_url, api_key=es_api_key, verify_certs=False, request_timeout=REQUEST_TIMEOUT
        )
        client.info()
        return client


def ensure_index(client: Elasticsearch) -> None:
    if client.indices.exists(index=INDEX_NAME):
        print(f"Index '{INDEX_NAME}' already exists, skipping creation")
        return
    client.indices.create(index=INDEX_NAME, body=INDEX_MAPPING)
    client.cluster.health(index=INDEX_NAME, wait_for_status="yellow", timeout="30s")
    print(f"Created index '{INDEX_NAME}'")


def to_float(value):
    value = (value or "").strip()
    if not value:
        return None
    try:
        return float(value)
    except ValueError:
        return None


def read_docs(csv_path: str):
    with open(csv_path, newline="", encoding="utf-8") as f:
        reader = csv.DictReader(f)
        for row in reader:
            doc = {
                "names": (row.get("names") or "").strip(),
                "date_x": (row.get("date_x") or "").strip(),
                "score": to_float(row.get("score")),
                "genre": [g.strip() for g in (row.get("genre") or "").split(",") if g.strip()],
                "overview": (row.get("overview") or "").strip(),
                "crew": (row.get("crew") or "").strip(),
                "orig_title": (row.get("orig_title") or "").strip(),
                "status": (row.get("status") or "").strip(),
                "orig_lang": (row.get("orig_lang") or "").strip(),
                "budget_x": to_float(row.get("budget_x")),
                "revenue": to_float(row.get("revenue")),
                "country": (row.get("country") or "").strip(),
            }
            yield {"_index": INDEX_NAME, "_source": doc}


def main() -> None:
    load_dotenv()
    es_url = os.environ["ES_URL"]
    es_api_key = os.environ["ES_API_KEY"]

    client = build_client(es_url, es_api_key)
    ensure_index(client)

    try:
        success, errors = bulk(
            client,
            read_docs(CSV_PATH),
            chunk_size=BULK_CHUNK_SIZE,
            raise_on_error=False,
        )
    except BulkIndexError as e:
        print(f"Bulk indexing failed: {e}", file=sys.stderr)
        sys.exit(1)

    print(f"Indexed {success} documents into '{INDEX_NAME}'")
    if errors:
        print(f"{len(errors)} documents failed to index", file=sys.stderr)
        for err in errors[:5]:
            print(err, file=sys.stderr)


if __name__ == "__main__":
    main()

运行上面的代码后,我们可以看到如下的结果:

GET imdb/_mapping

使用 Agent Builder 对数据进行探索

我们首先使用默认的 Elastic AI Agent 来对数据进行探索:

我们选择自己的大模型。如果大家还不知道如何创建连接器,请观看我的视频: 为 AI 安装自己的 Elasticsearch 集群

我们可以看出来大模型可以为我们生成想要的统计结果。我们再来问一下:

哪一年的生产的电影最多?

我们可以看到所使用的 ES|QL 查询。

电影的类别有哪些?那个类别的电影最多?

是不是使用默认的 Elastic AI Agent 就可以做很多的事了?

你如果想定制自己的 Agents,那么你可以参考我之前的系列文章:

  • Elastic AI agent builder 介绍()()()()(

最后祝大家学习愉快!

更多推荐