Embedding Models

qianmoQqianmoQ· 更新于 2026-10-08· 阅读 53 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

Embedding Models

This page covers text-based embedding models. Flink Agents does not currently support multimodal embeddings.

Overview

Embedding models convert text strings into high-dimensional vectors that capture semantic meaning, enabling powerful semantic search and retrieval capabilities. These vector representations allow agents to understand and work with text similarity, semantic search, and knowledge retrieval patterns.

In Flink Agents, embedding models are essential for:

  • Semantic Search: Finding relevant documents or information based on meaning rather than exact keyword matches
  • Text Similarity: Measuring how similar two pieces of text are in meaning
  • Knowledge Retrieval: Enabling agents to find and retrieve relevant context from large knowledge bases
  • Vector Databases: Storing and querying embeddings for efficient similarity search

Getting Started

To use embedding models in your agents, you need to define both a connection and setup using decorators/annotations, then access the embedding model through the runtime context.

Resource Declaration

Flink Agents provides decorators(in python) and annotations(in java) to simplify embedding model setup within agents:

Declare an embedding model connection

The @embedding_model_connection decorator/ @EmbeddingModelConnection annotation marks a method that creates an embedding model connection. This is typically defined once and shared across multiple embedding model setups.

Python

@embedding_model_connection
@staticmethod
def embedding_model_connection() -> ResourceDescriptor:
    ...

Java

@EmbeddingModelConnection
public static ResourceDescriptor embeddingModelConnection() {
    ...
}

Declare an embedding model setup

The @embedding_model_setup decorator/ @EmbeddingModelSetup annotation marks a method that creates an embedding model setup. This references an embedding model connection and adds embed-specific configuration like model and dimensions.

Python

@embedding_model_setup
@staticmethod
def embedding_model_setup() -> ResourceDescriptor:
    ...

Java

@EmbeddingModelSetup
public static ResourceDescriptor embeddingModelSetup() {
    ...
}

Usage Example

Here’s how to define and use embedding models in your agent:

Python

class MyAgent(Agent):

    @embedding_model_connection
    @staticmethod
    def openai_connection() -> ResourceDescriptor:
        return ResourceDescriptor(
            clazz=ResourceName.EmbeddingModel.OPENAI_CONNECTION,
            api_key="your-api-key-here",
            base_url="https://api.openai.com/v1",
            request_timeout=30.0
        )

    @embedding_model_setup
    @staticmethod
    def openai_embedding() -> ResourceDescriptor:
        return ResourceDescriptor(
            clazz=ResourceName.EmbeddingModel.OPENAI_SETUP,
            connection="openai_connection",
            model="your-embedding-model-here"
        )

    @action(InputEvent.EVENT_TYPE)
    @staticmethod
    def process_text(event: Event, ctx: RunnerContext) -> None:
        # Get the embedding model from the runtime context
        embedding_model = ctx.get_resource("openai_embedding", ResourceType.EMBEDDING_MODEL)

        # Use the embedding model to generate embeddings
        input_event = InputEvent.from_event(event)
        user_query = str(input_event.input)
        embedding = embedding_model.embed(user_query)

        # Handle the embedding
        # Process the embedding vector as needed for your use case

Java

public class MyAgent extends Agent {

    @EmbeddingModelConnection
    public static ResourceDescriptor ollamaConnection() {
        return ResourceDescriptor.Builder.newBuilder(
                        ResourceName.EmbeddingModel.OLLAMA_CONNECTION)
                .addInitialArgument("host", "http://localhost:11434")
                .build();
    }

    @EmbeddingModelSetup
    public static ResourceDescriptor ollamaEmbedding() {
        return ResourceDescriptor.Builder.newBuilder(ResourceName.EmbeddingModel.OLLAMA_SETUP)
                .addInitialArgument("connection", "ollamaConnection")
                .addInitialArgument("model", "nomic-embed-text")
                .build();
    }

    @Action(listenEventTypes = {InputEvent.EVENT_TYPE})
    public static void processText(Event event, RunnerContext ctx)
            throws Exception {
        InputEvent inputEvent = InputEvent.fromEvent(event);
        // Get the embedding model from the runtime context
        BaseEmbeddingModelSetup embeddingModel =
                (BaseEmbeddingModelSetup)
                        ctx.getResource("embeddingModel", ResourceType.EMBEDDING_MODEL);

        // Use the embedding model to generate embeddings
        String input = (String) inputEvent.getInput();
        float[] embedding = embeddingModel.embed(input);

        // Handle the embedding
        // Process the embedding vector as needed for your use case
    }
}

Built-in Providers

Amazon Bedrock

Amazon Bedrock provides embedding capabilities through the Amazon Titan Text Embeddings V2 model via the InvokeModel API. The integration supports configurable output dimensions (256, 512, or 1024) and parallelizes batch embedding via a configurable thread pool, since the Titan V2 model processes one text per API call. Authentication is handled via SigV4 using the AWS default credentials chain.

Amazon Bedrock embedding models are only supported in Java currently. To use Amazon Bedrock embeddings from Python agents, see Using Cross-Language Providers.

Prerequisites

  1. An AWS account with Amazon Bedrock model access enabled for Amazon Titan Text Embeddings V2
  2. IAM credentials configured via any method supported by the AWS Default Credentials Provider

BedrockEmbeddingModelConnection Parameters

Java

ParameterTypeDefaultDescription
regionString"us-east-1"AWS region for the Bedrock service
modelString"amazon.titan-embed-text-v2:0"Default embedding model ID
embed_concurrencyint4Thread pool size for parallel batch embedding
max_retriesint5Maximum number of API retry attempts (retries on throttling, 429, 503)

BedrockEmbeddingModelSetup Parameters

Java

ParameterTypeDefaultDescription
connectionStringRequiredReference to connection method name
modelStringNoneOverride the default embedding model from the connection
dimensionsintNoneOutput embedding dimensions: 256, 512, or 1024

Usage Example

Java

public class MyAgent extends Agent {

    @EmbeddingModelConnection
    public static ResourceDescriptor bedrockEmbeddingConnection() {
        return ResourceDescriptor.Builder.newBuilder(ResourceName.EmbeddingModel.BEDROCK_CONNECTION)
                .addInitialArgument("region", "us-east-1")
                .addInitialArgument("embed_concurrency", 8)
                .build();
    }

    @EmbeddingModelSetup
    public static ResourceDescriptor bedrockEmbedding() {
        return ResourceDescriptor.Builder.newBuilder(ResourceName.EmbeddingModel.BEDROCK_SETUP)
                .addInitialArgument("connection", "bedrockEmbeddingConnection")
                .addInitialArgument("model", "amazon.titan-embed-text-v2:0")
                .addInitialArgument("dimensions", 1024)
                .build();
    }

    ...
}

Available Models

The Bedrock embedding integration currently supports:

  • Amazon Titan Text Embeddings V2 (amazon.titan-embed-text-v2:0): supports 256, 512, or 1024 dimensions

The integration always requests normalized embeddings (unit vectors), which makes cosine similarity equivalent to dot product. If you need raw, un-normalized vectors, use a custom provider.

Visit the Amazon Bedrock Embedding Models documentation for the latest information.

Model availability varies by AWS region and requires explicit model access enablement in the Bedrock console. Always check the Amazon Bedrock documentation for regional availability before implementing in production.

Ollama

Ollama provides local embedding models that run on your machine, offering privacy and control over your data.

Prerequisites

  1. Install Ollama from https://ollama.com/
  2. Start the Ollama server: ollama serve
  3. Download an embedding model: ollama pull nomic-embed-text

OllamaEmbeddingModelConnection Parameters

Python

ParameterTypeDefaultDescription
base_urlstr"http://localhost:11434"Ollama server URL
request_timeoutfloat30.0HTTP request timeout in seconds

Java

ParameterTypeDefaultDescription
hostString"http://localhost:11434"Ollama server URL
modelStringnomic-embed-textName of the default embedding model

OllamaEmbeddingModelSetup Parameters

Python

ParameterTypeDefaultDescription
connectionstrRequiredReference to connection method name
modelstrRequiredName of the embedding model to use
truncateboolTrueWhether to truncate text exceeding model limits
keep_alivestr/float"5m"How long to keep model loaded in memory
additional_kwargsdict{}Additional Ollama API parameters

Java

ParameterTypeDefaultDescription
connectionStringRequiredReference to connection method name
modelStringRequiredName of the embedding model to use

Usage Example

Python

class MyAgent(Agent):

    @embedding_model_connection
    @staticmethod
    def ollama_connection() -> ResourceDescriptor:
        return ResourceDescriptor(
            clazz=ResourceName.EmbeddingModel.OLLAMA_CONNECTION,
            base_url="http://localhost:11434",
            request_timeout=30.0
        )

    @embedding_model_setup
    @staticmethod
    def ollama_embedding() -> ResourceDescriptor:
        return ResourceDescriptor(
            clazz=ResourceName.EmbeddingModel.OLLAMA_SETUP,
            connection="ollama_connection",
            model="nomic-embed-text",
            truncate=True,
            keep_alive="5m"
        )

    ...

Java

public class MyAgent extends Agent {

    @EmbeddingModelConnection
    public static ResourceDescriptor ollamaConnection() {
        return ResourceDescriptor.Builder.newBuilder(
                        ResourceName.EmbeddingModel.OLLAMA_CONNECTION)
                .addInitialArgument("host", "http://localhost:11434")
                .build();
    }

    @EmbeddingModelSetup
    public static ResourceDescriptor ollamaEmbedding() {
        return ResourceDescriptor.Builder.newBuilder(ResourceName.EmbeddingModel.OLLAMA_SETUP)
                .addInitialArgument("connection", "ollamaConnection")
                .addInitialArgument("model", "nomic-embed-text")
                .build();
    }

    ...
}

Available Models

Visit the Ollama Embedding Models Library for the complete and up-to-date list of available embedding models.

Some popular options include:

  • nomic-embed-text
  • all-minilm
  • mxbai-embed-large

Model availability and specifications may change. Always check the official Ollama documentation for the latest information before implementing in production.

OpenAI

OpenAI provides cloud-based embedding models with state-of-the-art performance.

OpenAI embedding models are currently supported in the Python API only. To use OpenAI from Java agents, see Using Cross-Language Providers.

Prerequisites

  1. Get an API key from OpenAI Platform

Usage Example

class MyAgent(Agent):

    @embedding_model_connection
    @staticmethod
    def openai_connection() -> ResourceDescriptor:
        return ResourceDescriptor(
            clazz=ResourceName.EmbeddingModel.OPENAI_CONNECTION,
            api_key="your-api-key-here",
            base_url="https://api.openai.com/v1",
            request_timeout=30.0,
            max_retries=3
        )

    @embedding_model_setup
    @staticmethod
    def openai_embedding() -> ResourceDescriptor:
        return ResourceDescriptor(
            clazz=ResourceName.EmbeddingModel.OPENAI_SETUP,
            connection="openai_connection",
            model="your-embedding-model-here",
            encoding_format="float"
        )

OpenAIEmbeddingModelConnection Parameters

ParameterTypeDefaultDescription
api_keystrRequiredOpenAI API key for authentication
base_urlstr"https://api.openai.com/v1"OpenAI API base URL
request_timeoutfloat30.0HTTP request timeout in seconds
max_retriesint3Maximum number of retry attempts
organizationstrNoneOptional organization ID
projectstrNoneOptional project ID

OpenAIEmbeddingModelSetup Parameters

ParameterTypeDefaultDescription
connectionstrRequiredReference to connection method name
modelstrRequiredOpenAI embedding model name
encoding_formatstr"float"Return format (“float” or “base64”)
dimensionsintNoneOutput dimensions (text-embedding-3 models only)
userstrNoneEnd-user identifier for monitoring
additional_kwargsdict{}Additional parameters for the OpenAI embeddings API

Available Models

Visit the OpenAI Embeddings documentation for the complete and up-to-date list of available embedding models.

Current popular models include:

  • text-embedding-3-small
  • text-embedding-3-large
  • text-embedding-ada-002

Model availability and specifications may change. Always check the official OpenAI documentation for the latest information before implementing in production.

Tongyi (DashScope)

Tongyi provides cloud-based embedding models from Alibaba Cloud, with strong support for Chinese and English text.

Tongyi embedding models are currently supported in the Python API only. To use Tongyi from Java agents, see Using Cross-Language Providers.

Prerequisites

  1. Get an API key from Alibaba Cloud DashScope

Usage Example

class MyAgent(Agent):

    @embedding_model_connection
    @staticmethod
    def tongyi_connection() -> ResourceDescriptor:
        return ResourceDescriptor(
            clazz=ResourceName.EmbeddingModel.TONGYI_CONNECTION,
            api_key="your-api-key-here",  # Or set DASHSCOPE_API_KEY env var
            request_timeout=30.0
        )

    @embedding_model_setup
    @staticmethod
    def tongyi_embedding() -> ResourceDescriptor:
        return ResourceDescriptor(
            clazz=ResourceName.EmbeddingModel.TONGYI_SETUP,
            connection="tongyi_connection",
            model="text-embedding-v4",
            text_type="query"
        )

TongyiEmbeddingModelConnection Parameters

ParameterTypeDefaultDescription
api_keystr$DASHSCOPE_API_KEYDashScope API key for authentication
request_timeoutfloat30.0HTTP request timeout in seconds

TongyiEmbeddingModelSetup Parameters

ParameterTypeDefaultDescription
connectionstrRequiredReference to connection method name
modelstr"text-embedding-v4"Embedding model name
text_typestrNoneInput type: "query" or "document"
dimensionintNoneOutput vector dimensions (model-dependent)
additional_kwargsdict{}Additional DashScope API parameters

Available Models

Visit the DashScope Embedding Models documentation for the complete and up-to-date list of available embedding models.

Some popular options include:

  • text-embedding-v4 (default, recommended)
  • text-embedding-v3
  • text-embedding-v2
  • text-embedding-v1

Model availability and specifications may change. Always check the official DashScope documentation for the latest information before implementing in production.

Using Cross-Language Providers

Flink Agents supports cross-language embedding model integration, allowing you to use embedding models implemented in one language (Java or Python) from agents written in the other language. This is particularly useful when an embedding model provider is only available in one language (e.g., OpenAI embedding is currently Python-only).

Limitations:

  • Cross-language resources are currently supported only when running in Flink, not in local development mode
  • Complex object serialization between languages may have limitations

How To Use

To leverage embedding model supports provided in a different language, you need to declare the resource within a built-in cross-language wrapper, and specify the target provider as an argument:

  • Using Java embedding models in Python: Use ResourceName.EmbeddingModel.JAVA_WRAPPER_CONNECTION and ResourceName.EmbeddingModel.JAVA_WRAPPER_SETUP, specifying the Java provider class via the java_clazz parameter
  • Using Python embedding models in Java: Use ResourceName.EmbeddingModel.PYTHON_WRAPPER_CONNECTION and ResourceName.EmbeddingModel.PYTHON_WRAPPER_SETUP, specifying the Python provider via the pythonClazz parameter

Usage Example

Using Java Embedding Model in Python

class MyAgent(Agent):

    @embedding_model_connection
    @staticmethod
    def java_embedding_connection() -> ResourceDescriptor:
        # In pure Java, the equivalent ResourceDescriptor would be:
        # ResourceDescriptor.Builder
        #     .newBuilder(ResourceName.EmbeddingModel.OLLAMA_CONNECTION)
        #     .addInitialArgument("host", "http://localhost:11434")
        #     .build();
        return ResourceDescriptor(
            clazz=ResourceName.EmbeddingModel.JAVA_WRAPPER_CONNECTION,
            java_clazz=ResourceName.EmbeddingModel.Java.OLLAMA_CONNECTION,
            host="http://localhost:11434"
        )

    @embedding_model_setup
    @staticmethod
    def java_embedding_model() -> ResourceDescriptor:
        # In pure Java, the equivalent ResourceDescriptor would be:
        # ResourceDescriptor.Builder
        #     .newBuilder(ResourceName.EmbeddingModel.OLLAMA_SETUP)
        #     .addInitialArgument("connection", "java_embedding_connection")
        #     .addInitialArgument("model", "nomic-embed-text")
        #     .build();
        return ResourceDescriptor(
            clazz=ResourceName.EmbeddingModel.JAVA_WRAPPER_SETUP,
            java_clazz=ResourceName.EmbeddingModel.Java.OLLAMA_SETUP,
            connection="java_embedding_connection",
            model="nomic-embed-text"
        )

    @action(InputEvent.EVENT_TYPE)
    @staticmethod
    def process_input(event: Event, ctx: RunnerContext) -> None:
        # Use the Java embedding model from Python
        input_event = InputEvent.from_event(event)
        embedding_model = ctx.get_resource("java_embedding_model", ResourceType.EMBEDDING_MODEL)
        embedding = embedding_model.embed(str(input_event.input))
        # Process the embedding vector as needed

Using Python Embedding Model in Java

public class MyAgent extends Agent {

    @EmbeddingModelConnection
    public static ResourceDescriptor pythonEmbeddingConnection() {
        // In pure Python, the equivalent ResourceDescriptor would be:
        // ResourceDescriptor(
        //     clazz=ResourceName.EmbeddingModel.OLLAMA_CONNECTION,
        //     base_url="http://localhost:11434"
        // )
        return ResourceDescriptor.Builder.newBuilder(ResourceName.EmbeddingModel.PYTHON_WRAPPER_CONNECTION)
                .addInitialArgument("pythonClazz", ResourceName.EmbeddingModel.Python.OLLAMA_CONNECTION)
                .addInitialArgument("base_url", "http://localhost:11434")
                .build();
    }

    @EmbeddingModelSetup
    public static ResourceDescriptor pythonEmbeddingModel() {
        // In pure Python, the equivalent ResourceDescriptor would be:
        // ResourceDescriptor(
        //     clazz=ResourceName.EmbeddingModel.OLLAMA_SETUP,
        //     connection="ollama_connection",
        //     model="nomic-embed-text"
        // )
        return ResourceDescriptor.Builder.newBuilder(ResourceName.EmbeddingModel.PYTHON_WRAPPER_SETUP)
                .addInitialArgument("pythonClazz", ResourceName.EmbeddingModel.Python.OLLAMA_SETUP)
                .addInitialArgument("connection", "pythonEmbeddingConnection")
                .addInitialArgument("model", "nomic-embed-text")
                .build();
    }

    @Action(listenEventTypes = {InputEvent.EVENT_TYPE})
    public static void processInput(Event event, RunnerContext ctx) throws Exception {
        InputEvent inputEvent = InputEvent.fromEvent(event);
        // Use the Python embedding model from Java
        BaseEmbeddingModelSetup embeddingModel =
            (BaseEmbeddingModelSetup) ctx.getResource(
                "pythonEmbeddingModel",
                ResourceType.EMBEDDING_MODEL);
        float[] embedding = embeddingModel.embed((String) inputEvent.getInput());
        // Process the embedding vector as needed
    }
}

Custom Providers

The custom provider APIs are experimental and unstable, subject to incompatible changes in future releases.

If you want to use embedding models not offered by the built-in providers, you can extend the base embedding classes and implement your own! The embedding system is built around two main abstract classes:

BaseEmbeddingModelConnection

Handles the connection to embedding services and provides the core embedding functionality.

Python

class MyEmbeddingConnection(BaseEmbeddingModelConnection):

    @abstractmethod
    def embed(self, text: str | Sequence[str], **kwargs: Any) -> list[float] | list[list[float]]:
        # Core method: convert text to embedding vector
        # - text: Input text to embed
        # - kwargs: Additional parameters from model_kwargs
        # - Returns: List of float values representing the embedding
        pass

Java

public class MyEmbeddingConnection extends BaseEmbeddingModelConnection {

    @Override
    public float[] embed(String text, Map<String, Object> parameters) {
        // Core method: convert text to embedding vector
        // - text: Input text to embed
        // - parameters: Additional parameters
        // - Returns: Float array representing the embedding
        float[] embedding = ...;
        return embedding;
    }

    @Override
    public List<float[]> embed(List<String> texts, Map<String, Object> parameters) {
        // Core method: convert texts to embedding vectors
        // - text: Input texts to embed
        // - parameters: Additional parameters
        // - Returns: List of float array representing the embeddings
        List<float[]> embeddings = ...;
        return embeddings;
    }
}

BaseEmbeddingModelSetup

The setup class acts as a high-level configuration interface that defines which connection to use and how to configure the embedding model.

Python

class MyEmbeddingSetup(BaseEmbeddingModelSetup):
    # Add your custom configuration fields here

    @property
    def model_kwargs(self) -> Dict[str, Any]:
        # Return model-specific configuration passed to embed()
        # This dictionary is passed as **kwargs to the embed() method
        return {"model": self.model, ...}

Java

public class MyEmbeddingSetup extends BaseEmbeddingModelSetup {

    @Override
    public Map<String, Object> getParameters() {
        // Return model-specific configuration passed to embed()
        // This dictionary is passed as parameters to the embed() method
        Map<String, Object> parameters = new HashMap<>();

        if (model != null) {
            parameters.put("model", model);
        }
        ...

        return parameters;
    }

}

评论

登录后参与评论

正在加载评论…