すべてのプロダクト
Search
ドキュメントセンター

OpenSearch:RAG に基づく対話型検索の構築

最終更新日:Jun 23, 2026

ナレッジベースに対する対話型検索のために、AI Search Open Platform は、検索拡張生成 (RAG) アプリケーションを構築するための完全な開発パイプラインを提供します。このパイプラインには、データ前処理、取得、回答生成の 3 つの主要モジュールが含まれています。AI Search Open Platform は、ドキュメント解析、再ランキング、回答生成など、各モジュールの機能をコンポーネントベースの構成可能なアルゴリズムサービスとして提供し、開発コードを迅速に生成できます。プラットフォームはこれらのサービスを API を通じて公開します。本トピックで説明するように、提供されたコードをダウンロードし、API キー、サービスエンドポイント、およびローカルのナレッジベース情報を置き換えることで、RAG に基づく対話型検索アプリケーションを迅速に構築できます。

仕組み

検索拡張生成 (RAG) は、情報検索と大規模言語モデル (LLM) を組み合わせ、生成されるコンテンツの関連性、精度、多様性を向上させる AI 技術です。クエリを処理する際、RAG システムはまず外部のナレッジベースから最も関連性の高い情報スニペットを取得します。次に、取得した情報と元のクエリをコンテキストとして LLM に提供します。このアプローチにより、モデルは内部のトレーニングデータのみに依存するのではなく、外部の最新データやドメイン固有のデータに基づいて応答を生成するため、より正確で情報に基づいた回答を生成できます。

RAG-based intelligent Q&A implementation flowchart.jpg

ユースケース

ナレッジベースに対する対話型検索は、社内ナレッジの取得や特定分野の専門的な Q&A などのユースケースに最適です。検索拡張生成 (RAG) と大規模言語モデル (LLM) を専門的なナレッジベースのドキュメントに適用することで、システムは複雑な自然言語クエリを理解し、応答できるようになります。これにより、エンタープライズユーザーは、PDF、Word、表、画像など、さまざまなドキュメント形式で必要な情報を迅速に見つけることができます。

対話型検索インターフェイスでは、ユーザーが自然言語で質問を送信します。システムは構造化された回答を返し、さらなる探索のために推奨されるフォローアップの質問を提供します。

前提条件

  • AI Search Open Platform サービスが有効化されていること。詳細については、「サービスの有効化」をご参照ください。

  • サービスエンドポイントと認証情報を取得していること。詳細については、「サービスエンドポイントの取得」および「API キーの管理」をご参照ください。

    AI Search Open Platform は、パブリックエンドポイントおよび VPC エンドポイント経由でのサービス呼び出しをサポートしています。クロスリージョン呼び出しは VPC エンドポイント経由でサポートされます。現在、中国 (上海)、中国 (杭州)、中国 (深セン)、中国 (北京)、中国 (張家口)、および中国 (青島) リージョンのユーザーは、VPC エンドポイントを通じて AI Search Open Platform サービスを呼び出すことができます。

    [API キー] ページの上部には、「API キーはサービス呼び出しの権限に使用されます。安全に保管してください。API キーが漏洩した場合は、直ちに無効にしてください。一度に最大 10 個の API キーを有効にできます。」というメッセージが表示されます。[アクセスドメイン] セクションには、[パブリック API ドメイン][プライベート API ドメイン] が表示され、どちらも HTTPS アクセスをサポートしています。

  • バージョン 8.5 以降の Alibaba Cloud Elasticsearch クラスターが作成されていること。詳細については、「Alibaba Cloud Elasticsearch クラスターの作成」をご参照ください。パブリックネットワークまたは VPC 経由でクラスターにアクセスする場合は、ご利用のデバイスの IP アドレスをクラスターの IP アドレスホワイトリストに追加する必要があります。詳細については、「Elasticsearch クラスターのパブリックまたはプライベート IP アドレスホワイトリストの設定」をご参照ください。

  • aiohttp 3.8.6 および elasticsearch 8.14 パッケージがインストールされた Python 3.7 以降の環境があること。

RAG 開発パイプラインの構築

説明

利便性のために、AI Search Open Platform は 4 種類の開発フレームワークを提供しています:

  • Java SDK

  • Python SDK

  • LangChain:ビジネスで既に LangChain フレームワークを使用している場合に選択します。

  • LlamaIndex:ビジネスで既に LlamaIndex フレームワークを使用している場合に選択します。

ステップ 1: サービスの選択とコードのダウンロード

ナレッジベースとビジネス要件に基づいて、RAG パイプラインのアルゴリズムサービスと開発フレームワークを選択します。このトピックでは、Python SDK を例として使用します。

  1. AI Search Open Platform コンソールにログインします。

  2. 中国 (上海) リージョンを選択し、AI 検索オープンプラットフォーム に切り替えてから、対象のワークスペースを選択します。

    説明
    • AI Search Open Platform は、中国 (上海) およびドイツ (フランクフルト) リージョンでのみ利用可能です。

    • 中国 (杭州)、中国 (深セン)、中国 (北京)、中国 (張家口)、および中国 (青島) リージョンのユーザーは、VPC エンドポイントを使用してクロスリージョンで AI Search Open Platform サービスを呼び出すことができます。

  3. 左側のナビゲーションウィンドウで、シナリオセンター をクリックします。RAG - ナレッジベースオンライン質疑応答 カードで、進む をクリックします。

  4. ドロップダウンリストから、ビジネス要件に基づいて必要なサービスを選択してください。各サービスの詳細情報は、詳細 タブで表示できます。

    説明
    • API を介して RAG パイプラインのアルゴリズムサービスを呼び出す場合、サービス ID (service_id) を提供する必要があります。例えば、ドキュメントコンテンツ解析サービスの ID は ops-document-analyze-001 です。

    • リストでサービスを切り替えると、生成されたコード内の service_id が自動的に更新されます。コードをダウンロードした後でも、service_id を変更して別のサービスを呼び出すことができます。

    段階

    サービスの説明

    ドキュメントコンテンツ解析

    ドキュメントコンテンツ解析サービス (ops-document-analyze-001):非構造化ドキュメント (テキスト、表、画像) からタイトルや段落などの論理構造を抽出し、構造化された形式で出力する汎用サービスです。

    画像コンテンツ解析

    • 画像コンテンツ理解サービス (ops-image-analyze-vlm-001):マルチモーダル大規模モデルを使用して、画像からテキストを解析、理解、認識します。抽出されたテキストは、画像検索や Q&A シナリオに使用できます。

    • 画像テキスト認識サービス (ops-image-analyze-ocr-001):OCR を使用して画像内のテキストを認識します。解析されたテキストは、画像検索や Q&A シナリオに使用できます。

    ドキュメントチャンキング

    ドキュメントチャンキングサービス (ops-document-split-001):HTML、Markdown、TXT 形式の構造化データを、ドキュメントの段落、テキストのセマンティクス、または指定されたルールに基づいて分割する汎用テキストチャンキングサービスです。また、ドキュメントからコード、画像、表をリッチテキストとして抽出することもサポートしています。

    テキスト埋め込み

    • OpenSearch テキスト埋め込みサービス-001 (ops-text-embedding-001):多言語 (40 以上) のテキスト埋め込みを提供します。最大入力長は 300 トークンで、出力ベクトル次元は 1536 です。

    • OpenSearch 汎用テキスト埋め込みサービス-002 (ops-text-embedding-002):多言語 (100 以上) のテキスト埋め込みを提供します。最大入力長は 8,192 トークンで、出力ベクトル次元は 1024 です。

    • OpenSearch テキスト埋め込みサービス-中国語-001 (ops-text-embedding-zh-001):中国語のテキスト埋め込みを提供します。最大入力長は 1,024 トークンで、出力ベクトル次元は 768 です。

    • OpenSearch テキスト埋め込みサービス-英語-001 (ops-text-embedding-en-001):英語のテキスト埋め込みを提供します。最大入力長は 512 トークンで、出力ベクトル次元は 768 です。

    疎テキスト埋め込み

    テキストデータを疎ベクトル表現に変換します。疎ベクトルはストレージ使用量が少なく、キーワードや用語の頻度を表すためによく使用されます。密ベクトルと組み合わせてハイブリッド検索を行い、取得性能を向上させることができます。

    OpenSearch 疎テキスト埋め込みサービス (ops-text-sparse-embedding-001):多言語 (100 以上) の疎テキスト埋め込みを提供します。最大入力長は 8,192 トークンです。

    クエリ分析

    クエリ分析サービス 001 (ops-query-analyze-001):大規模言語モデルを使用してユーザーのクエリ意図を理解し、類似の質問で拡張する汎用クエリ分析サービスです。

    検索エンジン

    • Alibaba Cloud Elasticsearch:オープンソースの Elasticsearch 上に構築されたフルマネージドのクラウドサービスです。オープンソースの機能と 100% 互換性があり、すぐに使える従量課金制のエクスペリエンスを提供します。

      説明

      検索エンジンとして Alibaba Cloud Elasticsearch を選択した場合、互換性の問題により疎テキスト埋め込みサービスは利用できません。代わりにテキスト埋め込みサービスを使用することを推奨します。

    • OpenSearch Vector Search Edition:Alibaba が開発した大規模な分散ベクトル検索エンジンです。さまざまなベクトル検索アルゴリズムをサポートし、高精度で優れたパフォーマンスを提供し、大規模でコスト効率の高いインデックス作成と取得を可能にします。そのインデックスは、水平スケーリングとマージ、ストリーミングビルド、リアルタイムクエリ、動的データ更新をサポートしています。

      説明

      OpenSearch Vector Search Edition を使用する必要がある場合は、RAG パイプラインのエンジン設定とコードを置き換えることができます。

    再ランキングサービス

    BGE reranker モデル (ops-bge-reranker-larger):汎用のドキュメントスコアリングサービスです。クエリとドキュメントコンテンツの関連性に基づいてドキュメントをソートし、スコアの高い順にランク付けして、スコアリング結果を出力します。

    大規模言語モデル

    • OpenSearch-Qwen-Turbo (ops-qwen-turbo):Qwen-Turbo 大規模言語モデル上に構築されたこのサービスは、教師あり学習でファインチューニングされ、検索拡張を強化し、有害な応答を削減します。

    • Qwen-Turbo (qwen-turbo):中国語や英語を含むさまざまな言語をサポートする Qwen シリーズの大規模言語モデルです。詳細については、「Qwen シリーズ LLM の概要」をご参照ください。

    • Qwen-Plus (qwen-plus):中国語や英語を含むさまざまな言語をサポートする Qwen-Turbo 大規模言語モデルの拡張版です。詳細については、「Qwen シリーズ LLM の概要」をご参照ください。

    • Qwen-Max (qwen-max):中国語や英語を含むさまざまな言語をサポートする Qwen シリーズの兆パラメータ級の超大規模言語モデルです。詳細については、「Qwen シリーズ LLM の概要」をご参照ください。

  5. サービスを選択した後、設定完了、コード照会へ をクリックしてコードを表示およびダウンロードします。

    コードは、RAG パイプラインのランタイムフローを反映する 2 つの部分、オフラインドキュメント処理とオンライン対話型検索に構造化されています。

    プロセス

    機能

    説明

    オフラインドキュメント処理

    解析、画像抽出、チャンキング、埋め込み、結果の Elasticsearch インデックスへの書き込みを含むドキュメント処理を行います。

    document_pipeline_execute メイン関数は、以下のワークフローを完了します。URL または Base64 エンコーディングを介してドキュメントを入力できます。

    1. ドキュメント解析。API の詳細については、「ドキュメント解析 API」をご参照ください。

      • 非同期ドキュメント解析 API を呼び出して、ドキュメント URL からコンテンツを抽出するか、Base64 エンコードされたファイルからコンテンツをデコードします。

      • create_async_extraction_task 関数を使用して解析タスクを作成し、poll_task_result 関数を使用してタスクの完了ステータスをポーリングします。

    2. 画像抽出。API の詳細については、「画像コンテンツ抽出 API」をご参照ください。

      • 非同期画像解析 API を呼び出して、画像 URL からコンテンツを抽出するか、Base64 エンコードされたファイルからデコードします。

      • create_image_analyze_task 関数を使用して画像解析タスクを作成し、get_image_analyze_task_status 関数を使用してそのステータスを取得します。

    3. ドキュメントチャンキング。API の詳細については、「ドキュメントチャンキング API」をご参照ください。

      • ドキュメントチャンキング API を呼び出して、指定された戦略に従って解析済みドキュメントを分割します。

      • ドキュメントチャンキングとリッチテキストコンテンツ解析の両方に document_split 関数を使用します。

    4. テキスト埋め込み。API の詳細については、「テキスト埋め込み API」をご参照ください。

      • テキスト埋め込み API を呼び出して、チャンク化されたテキストのベクトル表現を作成します。

      • text_embedding 関数を使用して、各チャンクの埋め込みベクトルを計算します。

    5. Elasticsearch への書き込み。サービスの詳細については、「Elasticsearch の k 最近傍 (kNN) 検索機能の使用」をご参照ください。

      • ベクトルフィールド embedding とドキュメントコンテンツフィールド content を指定する Elasticsearch インデックス設定を作成します。

        重要

        Elasticsearch インデックスを作成すると、同じ名前の既存のインデックスは削除されます。偶発的なデータ損失を避けるために、コード内のインデックス名を変更してください。

      • helpers.async_bulk 関数を使用して、ベクトル化された結果を Elasticsearch インデックスに一括書き込みします。

    オンライン対話型検索

    クエリベクトルの生成、クエリ分析の実行、関連するドキュメントチャンクの取得、検索結果の再ランキング、最終的な回答の生成など、オンラインのユーザークエリを処理します。

    query_pipeline_execute メイン関数は、以下のワークフローを完了してユーザークエリを処理し、回答を返します。

    1. クエリのベクトル化。API の詳細については、「テキスト埋め込み API」をご参照ください。

      • テキスト埋め込み API を呼び出して、ユーザークエリをベクトルに変換します。

      • text_embedding 関数を使用してクエリベクトルを生成します。

    2. クエリ分析サービスの呼び出し。詳細については、「クエリ分析 API」をご参照ください。

      このサービスは、会話履歴を分析してユーザーの意図を特定し、類似の質問を生成します。

    3. 埋め込みチャンクの検索。サービスの詳細については、「Elasticsearch の k 最近傍 (kNN) 検索機能の使用」をご参照ください。

      • Elasticsearch を使用して、クエリベクトルに類似したドキュメントチャンクをインデックスから取得します。

      • AsyncElasticsearchsearch API と kNN クエリを組み合わせて類似検索を実行します。

    4. 再ランキングサービスの呼び出し。詳細については、「再ランキング API」をご参照ください。

      • 再ランキングサービス API を呼び出して、取得したチャンクをスコアリングおよびソートします。

      • documents_ranking 関数を使用して、ユーザークエリに基づいてドキュメントをスコアリングおよびソートします。

    5. 大規模言語モデルによる回答の生成。API の詳細については、「回答生成 API」をご参照ください。

      LLM サービスを呼び出し、llm_call 関数に取得結果とユーザークエリを使用して最終的な回答を生成します。

    コード照会 で、ドキュメント処理オンライン質疑応答 を選択し、次に コードのコピー または ファイルのダウンロード をクリックしてコードをローカルに保存します。

ステップ 2: パイプラインの設定とテスト

コードを offline.pyonline.py などの 2 つのローカルファイルにダウンロードした後、コード内の主要なパラメーターを設定する必要があります。

カテゴリ

パラメーター

説明

AI Search Open Platform

api_key

認証用の API キー。詳細については、「API キーの管理」をご参照ください。

aisearch_endpoint

API 呼び出し用のサービスエンドポイント。詳細については、「サービスエンドポイントの取得」をご参照ください。

説明

エンドポイント URL から "http://" プレフィックスを削除してください。

API 呼び出しは、パブリックエンドポイントと VPC エンドポイントの両方でサポートされています。

workspace_name

AI Search Open Platform 上のワークスペースの名前。

service_id

サービス ID。便宜上、service_id_config ディクショナリを使用して、offline.pyonline.py の両方のファイルで異なるサービスのサービス ID を設定できます。

# AI Search Open Platform の設定
api_key = "xxx"
host = "http://xxx.platform-cn-shanghai.opensearch.aliyuncs.com"
workspace_name = "default"
# サービス ID の設定
service_id_config = {"extract": "ops-document-analyze-001", "split": "ops-document-split-001", "emb": "ops-text-embedding-001"}

Elasticsearch 検索エンジン

es_host

Elasticsearch クラスターのエンドポイント。パブリックネットワークまたは VPC 経由でクラスターにアクセスする場合は、ご利用のデバイスの IP アドレスをクラスターの IP アドレスホワイトリストに追加する必要があります。詳細については、「Elasticsearch クラスターのパブリックまたはプライベート IP アドレスホワイトリストの設定」をご参照ください。

es_auth

Elasticsearch クラスターにアクセスするためのユーザー名とパスワード。ユーザー名は elastic で、パスワードはクラスター作成時に設定したものです。パスワードを忘れた場合は、リセットできます。詳細については、「インスタンスのアクセスパスワードのリセット」をご参照ください。

その他のパラメーター

サンプルデータを使用する場合、変更は不要です。

パラメーターを設定した後、Python 3.7 以降の環境で、まず offline.py スクリプトを実行し、次に online.py スクリプトを実行して結果をテストします。

ナレッジベースのドキュメントが「AI Search Open Platform の概要」である場合、次の質問をします:AI Search Open Platform は何ができますか?

次の出力が表示されます:

  • オフラインドキュメント処理の結果

    image analyze :https://img.alicdn.com/imgextra/i2/O1CN01bYc1m81RrcSAyOjMu_!!6000000002165-54-tps-60-60.apng
        https://img.alicdn.com/imgextra/i2/O1CN01bYc1m81RrcSAyOjMu_!!6000000002165-54-tps-60-60.apng is not analyzable.
        image analyze :https://help-static-aliyun-doc.aliyuncs.com/assets/img/zh-CN/3873436171/p802381.png
        image analyze :https://help-static-aliyun-doc.aliyuncs.com/assets/img/zh-CN/0650850271/p819277.png
        image analyze :https://help-static-aliyun-doc.aliyuncs.com/assets/img/zh-CN/0650850271/p819277.png
        image analyze ://gw.alicdn.com/tfs/TB1GxwdSXXXXXa.aXXXXXXXXXXX-65-70.gif
            https://gw.alicdn.com/tfs/TB1GxwdSXXXXXa.aXXXXXXXXXXX-65-70.gif is not analyzable.
        image analyze ://img.alicdn.com/tfs/TB1..50QpXXXX7XpXXXXXXXXXX-40-40.png
        image analyze :https://img.alicdn.com/tfs/TB1A0dINW6qK1RjSZFmXXX0PFXa-258-258.jpg
        image analyze ://gw.alicdn.com/tfs/TB1GxwdSXXXXXa.aXXXXXXXXXXX-65-70.gif
            https://gw.alicdn.com/tfs/TB1GxwdSXXXXXa.aXXXXXXXXXXX-65-70.gif is not analyzable.
        image analyze ://img.alicdn.com/tfs/TB1..50QpXXXX7XpXXXXXXXXXX-40-40.png
        image analyze ://gw.alicdn.com/tfs/TB1GxwdSXXXXXa.aXXXXXXXXXXX-65-70.gif
            https://gw.alicdn.com/tfs/TB1GxwdSXXXXXa.aXXXXXXXXXXX-65-70.gif is not analyzable.
        image analyze ://img.alicdn.com/tfs/TB1..50QpXXXX7XpXXXXXXXXXX-40-40.png
    text-embedding done
    OS write response:  {"status":"OK","code":200}
  • オンライン対話型検索の結果

    /opt/miniconda3/envs/QA-pytest-base-lib1/bin/python /Users/liu/codeRepos/QA-pytest-base-lib/rag/case/SDK/python_sdk_es_zx.py
    query analysis rewrite result:What can the OpenSearch AI Search Open Platform do?
    Final answer from the large model:  AI Search Open Platform provides intelligent search services that power the core search functions for Alibaba's businesses, including Taobao and Tmall, and offers intelligent search solutions to external clients across various industries. It features industry-specific query semantic understanding and machine learning ranking algorithms to help developers build high-quality intelligent search services.
    The platform is suitable for a wide range of scenarios, including but not limited to:
    - E-commerce and retail intelligent search
    - Content and news search
    - Gaming industry search
    - Healthcare industry search
    - Financial industry search
    AI Search Open Platform focuses on intelligent search and Retrieval-Augmented Generation (RAG) scenarios, providing component-based services and flexible calling mechanisms. It has built-in services for document parsing, document chunking, text embedding, retrieval, reranking, and large language models, enabling a one-stop, flexible development experience for AI search applications.
    Process finished with exit code 0
  • ソースコードファイル

    offline.py
    # RAG オフラインパイプライン - Elasticsearch エンジン
    # 環境要件:
    # Python 3.7 以降
    # Elasticsearch クラスター 8.5 以降。Alibaba Cloud Elasticsearch を使用する場合は、事前にサービスを有効化し、IP アドレスホワイトリストを設定する必要があります。https://www.alibabacloud.com/help/elasticsearch/latest/configure-a-public-or-private-ip-address-whitelist-for-an-elasticsearch-cluster をご参照ください。
    
    # パッケージ要件:
    # pip install alibabacloud_searchplat20240529
    # pip install elasticsearch
    
    # AI Search Open Platform の設定
    aisearch_endpoint = "xxx.platform-cn-shanghai.opensearch.aliyuncs.com"
    api_key = "OS-xxx"
    workspace_name = "default"
    service_id_config = {"extract": "ops-document-analyze-001",
                         "split": "ops-document-split-001",
                         "text_embedding": "ops-text-embedding-001",
                         "text_sparse_embedding": "ops-text-sparse-embedding-001",
                         "image_analyze": "ops-image-analyze-ocr-001"}
    
    # Elasticsearch の設定
    es_host = 'http://es-cn-xxx.public.elasticsearch.aliyuncs.com:9200'
    es_auth = ('elastic', 'xxx')
    
    # 入力ドキュメントの URL。この例では、AI Search Open Platform の製品ドキュメントを使用します。
    document_url = "https://www.alibabacloud.com/help/open-search/search-platform/product-overview/introduction-to-search-platform"
    
    import asyncio
    from typing import List
    from elasticsearch import AsyncElasticsearch
    from elasticsearch import helpers
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_searchplat20240529.client import Client
    from alibabacloud_searchplat20240529.models import GetDocumentSplitRequest, CreateDocumentAnalyzeTaskRequest, \
        CreateDocumentAnalyzeTaskRequestDocument, GetDocumentAnalyzeTaskStatusRequest, \
        GetDocumentSplitRequestDocument, GetTextEmbeddingRequest, GetTextEmbeddingResponseBodyResultEmbeddings, \
        GetTextSparseEmbeddingRequest, GetTextSparseEmbeddingResponseBodyResultSparseEmbeddings, \
        CreateImageAnalyzeTaskRequestDocument, CreateImageAnalyzeTaskRequest, CreateImageAnalyzeTaskResponse, \
        GetImageAnalyzeTaskStatusRequest, GetImageAnalyzeTaskStatusResponse
    
    
    async def poll_task_result(ops_client, task_id, service_id, interval=5):
        while True:
            request = GetDocumentAnalyzeTaskStatusRequest(task_id=task_id)
            response = await ops_client.get_document_analyze_task_status_async(workspace_name, service_id, request)
            status = response.body.result.status
            if status == "PENDING":
                await asyncio.sleep(interval)
            elif status == "SUCCESS":
                return response
            else:
                raise Exception("document analyze task failed")
    
    
    def is_analyzable_url(url:str):
        if not url:
            return False
        image_extensions = {'.jpg', '.jpeg', '.png', '.bmp', '.tiff'}
        return url.lower().endswith(tuple(image_extensions))
    
    
    async def image_analyze(ops_client, url):
        try:
            print("image analyze :" + url)
            if url.startswith("//"):
                url = "https:" + url
            if not is_analyzable_url(url):
                print(url + " is not analyzable.")
                return url
            image_analyze_service_id = service_id_config["image_analyze"]
            document = CreateImageAnalyzeTaskRequestDocument(
                url=url,
            )
            request = CreateImageAnalyzeTaskRequest(document=document)
            response: CreateImageAnalyzeTaskResponse = ops_client.create_image_analyze_task(workspace_name, image_analyze_service_id, request)
            task_id = response.body.result.task_id
            while True:
                request = GetImageAnalyzeTaskStatusRequest(task_id=task_id)
                response: GetImageAnalyzeTaskStatusResponse = ops_client.get_image_analyze_task_status(workspace_name, image_analyze_service_id, request)
                status = response.body.result.status
                if status == "PENDING":
                    await asyncio.sleep(5)
                elif status == "SUCCESS":
                    return url + response.body.result.data.content
                else:
                    print("image analyze error: " + response.body.result.error)
                    return url
        except Exception as e:
            print(f"image analyze Exception : {e}")
    
    
    def chunk_list(lst, chunk_size):
        for i in range(0, len(lst), chunk_size):
            yield lst[i:i + chunk_size]
    
    
    async def write_to_es(doc_list):
        es = AsyncElasticsearch(
            [es_host],
            basic_auth=es_auth,
            verify_certs=False,  # SSL 証明書を検証しない
            request_timeout=30,
            max_retries=10,
            retry_on_timeout=True
        )
        index_name = 'dense_vertex_index'
    
        # 既存のインデックスが存在する場合は削除します。
        if await es.indices.exists(index=index_name):
            await es.indices.delete(index=index_name)
    
        # ベクトルインデックスを作成します。`emb` フィールドを dense_vector、`content` を text、`source_doc` を keyword として指定します。
        index_mappings = {
            "mappings": {
                "properties": {
                    "emb": {
                        "type": "dense_vector",
                        "index": True,
                        "similarity": "cosine",
                        "dims": 1536  # テキスト埋め込みモデルの出力に基づいて次元を変更します。
                    },
                    "content": {
                        "type": "text"
                    },
                    "source_doc": {
                        "type": "keyword"
                    }
                }
            }
        }
        await es.indices.create(index=index_name, body=index_mappings)
    
        # 埋め込み結果を新しく作成したインデックスに一括アップロードします。
        actions = []
        for i, doc in enumerate(doc_list):
            action = {
                "_index": index_name,
                "_id": doc['id'],
                "_source": {
                    "emb": doc['embedding'],
                    "content": doc['content'],
                    "source_doc": document_url
                }
            }
            actions.append(action)
    
        try:
            await helpers.async_bulk(es, actions)
        except Exception as e:
            for error in e.errors:
                print(error)
    
        # アップロードが成功したことを確認します。
        await asyncio.sleep(2)
        query = {
            "query": {
                "ids": {
                    "values": [doc_list[0]["id"]]
                }
            }
        }
        res = await es.search(index=index_name, body=query)
        if len(res['hits']['hits']) > 0:
            print("ES write success")
        await es.close()
    
    
    async def document_pipeline_execute(document_url: str = None, document_base64: str = None, file_name: str = None):
    
        # AI Search Open Platform クライアントを初期化します。
        config = Config(bearer_token=api_key, endpoint=aisearch_endpoint, protocol="http")
        ops_client = Client(config=config)
    
        # ステップ 1: ドキュメント解析
        document_analyze_request = CreateDocumentAnalyzeTaskRequest(
            document=CreateDocumentAnalyzeTaskRequestDocument(url=document_url, content=document_base64,
                                                              file_name=file_name, file_type='html'))
        document_analyze_response = await ops_client.create_document_analyze_task_async(workspace_name=workspace_name,
                                                                                        service_id=service_id_config[
                                                                                            "extract"],
                                                                                        request=document_analyze_request)
        print("document-analyze task_id:" + document_analyze_response.body.result.task_id)
        extraction_result = await poll_task_result(ops_client, document_analyze_response.body.result.task_id,
                                                   service_id_config["extract"])
        print("document-analyze done")
        document_content = extraction_result.body.result.data.content
        content_type = extraction_result.body.result.data.content_type
        
        # ステップ 2: ドキュメントチャンキング
        document_split_request = GetDocumentSplitRequest(
            GetDocumentSplitRequestDocument(content=document_content, content_type=content_type))
        document_split_result = await ops_client.get_document_split_async(workspace_name, service_id_config["split"],
                                                                          document_split_request)
        print("document-split done, chunks count: " + str(len(document_split_result.body.result.chunks))
              + " rich text count:" + str(len(document_split_result.body.result.rich_texts)))
    
        # ステップ 3: テキスト埋め込み
        # チャンキング結果を抽出します。画像チャンクの場合、画像分析サービスを使用してテキストコンテンツを抽出します。
        doc_list = ([{"id": chunk.meta.get("id"), "content": chunk.content} for chunk in
                     document_split_result.body.result.chunks]
                    + [{"id": chunk.meta.get("id"), "content": chunk.content} for chunk in
                       document_split_result.body.result.rich_texts if chunk.meta.get("type") != "image"]
                    + [{"id": chunk.meta.get("id"), "content": await image_analyze(ops_client, chunk.content)} for chunk in
                       document_split_result.body.result.rich_texts if chunk.meta.get("type") == "image"]
                    )
    
        chunk_size = 32  # 1 リクエストあたり最大 32 個の埋め込みを計算できます。
        all_text_embeddings: List[GetTextEmbeddingResponseBodyResultEmbeddings] = []
        for chunk in chunk_list([text["content"] for text in doc_list], chunk_size):
            response = await ops_client.get_text_embedding_async(workspace_name, service_id_config["text_embedding"],
                                                                 GetTextEmbeddingRequest(chunk))
            all_text_embeddings.extend(response.body.result.embeddings)
    
        all_text_sparse_embeddings: List[GetTextSparseEmbeddingResponseBodyResultSparseEmbeddings] = []
        for chunk in chunk_list([text["content"] for text in doc_list], chunk_size):
            response = await ops_client.get_text_sparse_embedding_async(workspace_name,
                                                                        service_id_config["text_sparse_embedding"],
                                                                        GetTextSparseEmbeddingRequest(chunk,
                                                                                                      input_type="document",
                                                                                                      return_token=True))
            all_text_sparse_embeddings.extend(response.body.result.sparse_embeddings)
    
        for i in range(len(doc_list)):
            doc_list[i]["embedding"] = all_text_embeddings[i].embedding
            doc_list[i]["sparse_embedding"] = all_text_sparse_embeddings[i].embedding
    
        print("text-embedding done")
    
        # ステップ 4: Elasticsearch ストレージエンジンへの書き込み
        await write_to_es(doc_list)
    
    
    if __name__ == "__main__":
        # 非同期タスクを実行します。
        #    import nest_asyncio # Jupyter ノートブックで実行する場合は、この 2 行のコメントを解除してください。
        #    nest_asyncio.apply() # Jupyter ノートブックで実行する場合は、この 2 行のコメントを解除してください。
        asyncio.run(document_pipeline_execute(document_url))
        # asyncio.run(document_pipeline_execute(document_base64="eHh4eHh4eHg...", file_name="attention.pdf")) # 別の呼び出し方法
              
    online.py
    # RAG オンラインパイプライン - Elasticsearch エンジン
    # 環境要件:
    # Python 3.7 以降
    # Elasticsearch クラスター 8.5 以降。Alibaba Cloud Elasticsearch を使用する場合は、事前にサービスを有効化し、IP アドレスホワイトリストを設定する必要があります。https://www.alibabacloud.com/help/elasticsearch/latest/configure-a-public-or-private-ip-address-whitelist-for-an-elasticsearch-cluster をご参照ください。
    
    # パッケージ要件:
    # pip install alibabacloud_searchplat20240529
    # pip install elasticsearch
    
    # AI Search Open Platform の設定
    api_key = "OS-xxx"
    aisearch_endpoint = "xxx.platform-cn-shanghai.opensearch.aliyuncs.com"
    workspace_name = "default"
    service_id_config = {
        "rank": "ops-bge-reranker-larger",
        "text_embedding": "ops-text-embedding-001",
        "text_sparse_embedding": "ops-text-sparse-embedding-001",
        "llm": "ops-qwen-turbo",
        "query_analyze": "ops-query-analyze-001"
    }
    
    # Elasticsearch の設定
    es_host = 'http://es-cn-xxx.public.elasticsearch.aliyuncs.com:9200'
    es_auth = ('elastic', 'xxx')
    
    # ユーザーのクエリ:
    user_query = "What can AI Search Open Platform do?"
    
    import asyncio
    from elasticsearch import AsyncElasticsearch
    from alibabacloud_tea_openapi.models import Config
    from alibabacloud_searchplat20240529.client import Client
    from alibabacloud_searchplat20240529.models import GetTextEmbeddingRequest,  \
        GetDocumentRankRequest, GetTextGenerationRequest, GetTextGenerationRequestMessages, \
        GetQueryAnalysisRequest
    
    # AI Search Open Platform クライアントを初期化します。
    config = Config(bearer_token=api_key, endpoint=aisearch_endpoint, protocol="http")
    ops_client = Client(config=config)
    
    
    async def es_retrieve(query):
        es = AsyncElasticsearch(
            [es_host],
            basic_auth=es_auth,
            verify_certs=False,
            request_timeout=30,
            max_retries=10,
            retry_on_timeout=True
        )
        index_name = 'dense_vertex_index'
        # クエリをベクトル化します。
        query_emb_result = await ops_client.get_text_embedding_async(workspace_name, service_id_config["text_embedding"],
                                                                     GetTextEmbeddingRequest(input=[query],
                                                                                             input_type="query"))
        query_emb = query_emb_result.body.result.embeddings[0].embedding
        query = {
            "field": "emb",
            "query_vector": query_emb,
            "k": 5,  # 返すドキュメントチャンクの数
            "num_candidates": 100  # HNSW 検索パラメーター (ef_search)
        }
    
        res = await es.search(index=index_name, knn=query)
        search_results = [item['_source']['content'] for item in res['hits']['hits']]
        await es.close()
        return search_results
    
    
    # オンライン対話型検索パイプライン。入力はユーザーの質問です。
    async def query_pipeline_execute():
    
        # ステップ 1: クエリ分析
        query_analyze_response = ops_client.get_query_analysis(workspace_name, service_id_config['query_analyze'],
                                                               GetQueryAnalysisRequest(query=user_query))
        print("query analysis rewrite result:" + query_analyze_response.body.result.query)
    
        # ステップ 2: ドキュメント取得
        all_query_results = []
        user_query_results = await es_retrieve(user_query)
        all_query_results.extend(user_query_results)
        rewrite_query_results = await es_retrieve(query_analyze_response.body.result.query)
        all_query_results.extend(rewrite_query_results)
        for extend_query in query_analyze_response.body.result.queries:
            extend_query_result = await es_retrieve(extend_query)
            all_query_results.extend(extend_query_result)
        # 取得したすべての結果から重複を削除します。
        remove_duplicate_results = list(set(all_query_results))
    
        # ステップ 3: 取得したドキュメントの再ランキング
        rerank_top_k = 8
        score_results = await ops_client.get_document_rank_async(workspace_name, service_id_config["rank"],GetDocumentRankRequest(remove_duplicate_results, user_query))
        rerank_results = [remove_duplicate_results[item.index] for item in score_results.body.result.scores[:rerank_top_k]]
    
        # ステップ 4: 大規模言語モデルを呼び出して回答を生成
        docs = '\n'.join([f"<article>{s}</article>" for s in rerank_results])
        messages = [
            GetTextGenerationRequestMessages(role="system", content="You are a helpful assistant."),
            GetTextGenerationRequestMessages(role="user",
                                             content=f"""The provided information contains multiple independent documents, each enclosed in <article> and </article> tags. Information:\n'''{docs}'''
                                             \n\nBased on the information provided above, answer the user's question in a detailed and organized manner. Ensure your answer fully addresses the question and correctly uses the provided information. If the information is insufficient to answer the question, state "The question cannot be answered based on the provided information." Do not use any information outside of the provided context. Ensure that every statement in your answer is supported by the context. Please answer in English.
                                             \nQuestion: '''{user_query}'''""""")
        ]
        response = await ops_client.get_text_generation_async(workspace_name, service_id_config["llm"],
                                                              GetTextGenerationRequest(messages=messages))
        print("Final answer from the large model: ", response.body.result.text)
    
    
    if __name__ == "__main__":
        # 非同期タスクを実行します。
        #    import nest_asyncio # Jupyter ノートブックで実行する場合は、この 2 行のコメントを解除してください。
        #    nest_asyncio.apply() # Jupyter ノートブックで実行する場合は、この 2 行のコメントを解除してください。
        asyncio.run(query_pipeline_execute())
              

よくある質問

コードの実行中に、リソースが時間内に解放されないために「Unclosed connector」というメッセージが表示されることがあります。このメッセージは無視して問題ありません。