Vector Elasticsearch Sink 完全指南:Bulk 与 Data Streams 写入、API 版本自动检测及源码级配置解析
【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector
本篇技术指南围绕 Vector 的 Elasticsearch sink 展开,覆盖其全部配置项(端点、认证、bulk/data_stream 双模式、批量与压缩、版本控制等)、写入工作机制(冲突处理、部分失败重试、多端点分发、健康检查),并结合 src/sinks/elasticsearch/ 的源码实现,讲解 API 版本自动检测、模板限定(confinement)等底层细节,帮助你在生产环境中正确、安全地把日志与指标写入 Elasticsearch 或 OpenSearch。
组件概览
Elasticsearch sink 是 Vector 中将可观测性事件(logs 与各类 metrics,不支持 traces)批量写入 Elasticsearch / OpenSearch 的输出组件。根据 website/cue/reference/components/sinks/elasticsearch.cue 中的组件元数据,其核心能力可归纳为:
- 投递语义:
at_least_once(至少一次),无状态(stateful: false),出口方式为批量(batch); - 数据输入:支持 logs 及 counter / gauge / histogram / summary / distribution / set 全类型 metrics;
- 端到端确认(acknowledgements):支持;
- 健康检查:启动时主动探测(
enabled: true); - 批量发送:默认
max_bytes = 10_000_000、timeout_secs = 1.0; - 压缩:默认
none,可选gzip(配置层还支持snappy/zlib/zstd); - 代理、自定义请求头、TLS:均支持,且 TLS 默认不启用(
enabled_default: false),可通过 scheme 推断; - 适用服务商:AWS、Azure、Elastic、GCP;
- 一个重要前提:Data streams 特性要求 Vector 配置
create作为bulk.action,该行为并非默认开启(data_stream 模式下源码会自动强制为create,见下文)。
组件文档页面 website/content/en/docs/reference/configuration/sinks/elasticsearch.md 本身是模板占位页,实际内容(参数说明、how it works)由 website/cue/reference/components/sinks/elasticsearch.cue 与 website/cue/reference/components/sinks/generated/elasticsearch.cue 两份 CUE 数据驱动生成;而 CUE 中的每条描述又与 src/sinks/elasticsearch/config.rs 中ElasticsearchConfig结构体的字段文档一一对应。
基础配置示例
一个最简可用配置(TOML 风格,与 website/cue/reference/components/sinks/elasticsearch.cue 中 OpenSearch 兼容章节的示例一致):
[sinks.elastic_logs] type = "elasticsearch" endpoints = ["http://10.24.32.122:9000"] bulk.index = "vector-%Y.%m.%d"一个覆盖常用选项的完整示例:
sinks: my_elastic: type: elasticsearch endpoints: - "https://user:password@example.com" # 可内嵌 Basic 凭据 api_version: auto # auto / v6 / v7 / v8 id_key: "id" # 事件字段映射到 _id pipeline: "my-ingest-pipeline" compression: gzip mode: bulk # bulk / data_stream bulk: action: "index" # index / create / update,支持模板 {{ action }} index: "application-{{ application_id }}-%Y-%m-%d" template_fallback_index: "test-index" # 模板渲染失败时的兜底索引 version: "{{ obj_version }}-%Y-%m-%d" # 可选,外部版本控制 version_type: external # internal / external / external_gte data_stream: # 仅 mode = data_stream 时生效 type: "logs" dataset: "{{ service }}" namespace: "{{ environment }}" auto_routing: true sync_fields: true batch: max_events: 1000 timeout_secs: 1.0 request_retry_partial: false auth: strategy: basic # basic / aws user: "${ELASTICSEARCH_USERNAME}" password: "${ELASTICSEARCH_PASSWORD}" query: X-Powered-By: "Vector" tls: enabled: true完整配置项说明
以下参数说明继承自 website/cue/reference/components/sinks/generated/elasticsearch.cue,默认值与校验规则均已在源码中核实。
端点与网络
| 选项 | 类型 / 默认值 | 说明 |
|---|---|---|
endpoint | string,已废弃 | 单端点写法,已废弃,应使用endpoints。源码在 src/sinks/elasticsearch/common.rs 的parse_many中检测到该选项会打印 DEPRECATION 警告 |
endpoints | array,默认[] | 端点列表,每个元素必须包含 HTTP scheme,可带主机名/IP/端口,也可内嵌 Basic 凭据(如https://user:password@example.com)。与endpoint互斥(required_one_of),二者必须恰好设置一个 |
tls | object,可选 | 标准 TLS 配置,支持证书与主机名校验;默认不启用 |
query | object,可选 | 追加到每个 HTTP 请求 query string 的自定义参数;值可以是单值,也可以是数组(多值 key,见下文 “Query params structure”) |
端点校验逻辑:HttpEndpoint在反序列化阶段即拒绝缺少 scheme 之外的非法 URI(缺省 scheme 会归一为 https)与无主机名端点,这在 src/sinks/elasticsearch/config.rs 的validate_rejects_endpoint_without_host/validate_rejects_non_http_endpoint测试中有明确断言;endpoint与endpoints同时出现、或两者都缺失,分别报 “mutually exclusive” / “Endpoints option must be specified” 错误。
索引与模式
| 选项 | 类型 / 默认值 | 说明 |
|---|---|---|
mode | string,默认bulk | bulk:使用 Bulk API 的indexaction 批量写入;data_stream:使用createaction,并遵循 Data Streams 语义(ECS 兼容,自动把事件的timestamp字段重命名为@timestamp)。bulk有normal别名(serde(alias = "normal")) |
doc_type | string,默认_doc | 仅对 Elasticsearch <= 6.x 有意义;7.0+ 已移除该概念 |
bulk.action | string 模板,默认index | Bulk API 动作,仅支持index/create/update;支持模板(如{{ action }}),源码中以UnconfinedTemplate承载 |
bulk.index | string 模板,默认vector-%Y.%m.%d | 目标索引名,支持日期占位符与字段插值 |
bulk.template_fallback_index | string,可选 | 当bulk.index模板无法渲染时写入的兜底索引 |
bulk.version | string 模板,可选 | 文档版本号;必须能解析为整数 |
bulk.version_type | string,默认internal | internal/external(含external_gt)/external_gte |
id_key | string,可选 | 指定事件字段名映射到 Elasticsearch 的_id。默认不设置(由 ES 自动生成);文档提醒自定义 ID 可能影响索引性能。外部版本控制(external/external_gte)必须配合id_key使用 |
pipeline | string,可选 | Ingest Pipeline 名称;源码会将其作为pipelinequery 参数拼进 bulk URI |
外部版本控制的一致性约束(在 src/sinks/elasticsearch/config.rs 的validate与 src/sinks/elasticsearch/common.rs 的parse_config中双重校验):
bulk.version已设置但version_type = internal→ 报错ExternalVersionIgnoredWithInternalVersioning;bulk.version已设置、version_type为 external 系但缺少id_key→ 报错ExternalVersioningWithoutDocumentID;bulk.version未设置但version_type为 external 系 → 报错ExternalVersioningWithoutVersion。
模板安全(confinement):bulk.index与data_stream.*的路由类模板默认处于“受限”状态——只有模板渲染出的索引名满足静态前缀等约定才允许通过,防止日志生产者通过字段值把事件路由到任意索引。该机制由ElasticsearchConfig::common_mode()(src/sinks/elasticsearch/config.rs)构建ConfinedTemplate实现;dangerously_allow_unconfined_template_resolution是一个显式的全局豁免开关,默认false,文档明确标注其“DANGEROUS — disables a security control”,会绕过启动校验与运行时限定。template_fallback_index仅在模板渲染失败(非安全类失败)时生效;遇到 confinement 违例时事件会被丢弃而不会落到兜底索引。
API 版本与 data stream 命名
| 选项 | 类型 / 默认值 | 说明 |
|---|---|---|
api_version | string,默认auto | auto自动探测;v6使用 6.x API;v7使用 7.x 兼容 API(含 OpenSearch);v8使用 8.x API。Amazon OpenSearch Serverless 必须保持auto |
suppress_type_name | bool,默认false,已废弃 | 是否发送type字段(7.x 废弃、8.x 移除),应改用api_version |
data_stream.type | string 模板,默认logs | data stream 名称三段式的 type 段 |
data_stream.dataset | string 模板,默认generic | dataset 段 |
data_stream.namespace | string 模板,默认default | namespace 段 |
data_stream.auto_routing | bool,默认true | 事件上存在data_stream.{type,dataset,namespace}字段时优先用事件字段推导 data stream 名(格式<type>-<dataset>-<namespace>),否则回落到配置值 |
data_stream.sync_fields | bool,默认true | 事件缺少data_stream.*字段时自动补齐并同步,保证字段与接收事件的 data stream 名一致 |
源码对应关系:DataStreamMode::index()按auto_routing分支组装名称(src/sinks/elasticsearch/config.rs);DataStreamMode::remap_timestamp()实现timestamp→@timestamp重命名(常量DATA_STREAM_TIMESTAMP_KEY)。值得注意的是,auto_routing取事件字段值前会经过两层校验:一是对应模板的 confinement 检查,二是is_valid_data_stream_component()的基础标识符校验(长度上限 100、禁止\ / * ? " < > | , # :空格等字符、禁止- _ + .开头、禁止./..段、dataset 与 namespace 额外禁止-),非法值触发TemplateRenderingError并丢弃事件。从源码注释看,这一层防御是为了让受攻击者控制的data_stream.*字段无法绕过构建期 confinement、也无法构造 Elasticsearch 本身会拒绝的索引名。
认证与 AWS
auth对象以strategy为标签字段,分basic与aws两种(src/sinks/elasticsearch/mod.rs 中ElasticsearchAuthConfig枚举,serde(tag = "strategy")):
- basic(HTTP Basic 认证):
user、password必填; - aws(Amazon OpenSearch Service 专用):
access_key_id、secret_access_key、assume_role(必填项,均可省略以走默认凭据链)、credentials_file(可选路径)、profile(默认default)、region(可选,缺省用服务自身区域)、external_id、session_name(缺省自动生成如assume-role-provider-<时间戳>)、session_token、imds(IMDS 配置对象)、load_timeout_secs(凭据加载超时,秒)。
关键约束(源码ParseError枚举与validate()实现):
- AWS 认证必须配置
aws.region,否则报aws.region required when AWS authentication is in use; - OpenSearch Serverless 必须
strategy = aws,否则报 “Amazon OpenSearch Serverless requiresauth.strategyvalue to beaws”; - Serverless 必须
api_version = auto,否则报 “Amazon OpenSearch Serverless requiresapi_versionvalue to beauto”。
opensearch_service_type取managed(默认,Elasticsearch 或托管 OpenSearch 域名)或serverless(OpenSearch Serverless 集合);源码 src/sinks/elasticsearch/config.rs 中OpenSearchServiceType::as_str()显示其最终决定 AWS v4 签名使用的服务名(es与aoss),并在 Serverless 下额外附带x-amz-content-sha256头参与签名(sign_request,src/sinks/elasticsearch/common.rs)。
批量、请求与容错
| 选项 | 类型 / 默认值 | 说明 |
|---|---|---|
batch | object | 标准批量行为(BatchConfig<RealtimeSizeBasedDefaultBatchSettings>),组件特性中默认max_bytes = 10MB、timeout_secs = 1.0 |
compression | string,默认none | none/gzip/snappy/zlib/zstd,未指定级别时一律使用默认压缩级别 |
encoding | Transformer | 序列化前对事件做转换(如metric_to_log相关行为) |
request | object | 出站 HTTP 设置(RequestConfig),含tower请求限制;超时会被自动写入 bulk URI 的timeoutquery 参数 |
request_retry_partial | bool,默认false | 是否对“整体成功但含部分失败”的 bulk 请求重试整个请求;官方建议同时使用id_key避免重复 |
acknowledgements | bool/object,默认false | 端到端确认行为(AcknowledgementsConfig) |
distribution | object,可选 | 端点健康判定选项(HealthConfig),配合多端点分发使用 |
bulk URI 组装:src/sinks/elasticsearch/common.rs 的parse_config将query自定义参数、自动追加的timeout=<请求超时>s、可选的pipeline=<名称>序列化为 query string,最终请求目标为{base_url}/_bulk?<params>。
metrics 输入支持
metrics对象(MetricToLogConfig)配置内置metric_to_log行为,将指标转为日志后再写入:
host_tag:指标上表示源主机的 tag 名,命中后写入生成日志的host字段(受全局log_schema.host_key影响);metric_tag_values:默认single(多值 tag 只展示最后一个非裸值);full把所有 tag 以数组形式展开;auto按底层形状编码(单值 tag 为字符串、多值 tag 为数组,长度 1 的数组往返为标量);timezone:对不含显式时区的时间戳转换所用时区,覆盖全局timezone选项,可用 TZ 库名称或local。
工作机制(How It Works)
以下各节完整继承组件 CUE 中how_it_works的官方说明,并补充源码印证。
冲突(Conflicts)
Vector 将数据批量缓冲后统一发送到 Elasticsearch 的_bulkAPI。默认所有事件以indexaction 写入:若已存在相同id的文档,会被替换。若把bulk.action配置为create,Elasticsearch 不会替换已存在文档,而是返回冲突错误。当bulk.action设为update时,文档在若干约束下被更新:message 必须放在.doc中且.doc_as_upsert为 true;update操作要求设置id_key,且encoding字段应指定doc与doc_as_upsert作为值。源码侧,BulkAction枚举(Index/Create/Update,src/sinks/elasticsearch/mod.rs)通过as_json_pointer()提供/index、/create、/update三种行头形式。
Data Streams
默认 Vector 使用indexaction 走 Bulk API。要使用 Elasticsearch Data Streams,需将mode设为data_stream,并改用data_stream.type/data_stream.dataset/data_stream.namespace的组合代替bulk.index。Data Streams 仅支持createaction——从源码结构看,ElasticsearchCommonMode::bulk_action()对DataStream分支直接返回Some(BulkAction::Create),绕过模板渲染,这是硬性保证而非配置。
多端点分发(Distribution)
当endpoints指定多个端点时,事件会按各端点的估计负载在其间分发,并支持故障转移(failover)。限流(rate limit)作用于整个 sink,而并发设置则逐端点管理。端点健康度被主动监控:由一个熔断器(circuit breaker)监控响应,连续失败达到阈值后触发,进入指数退避循环,每轮只放行单个请求试探该端点,收到成功响应后熔断器复位。构建阶段可见其装配:validated.request_limits.distributed_service(ElasticsearchRetryLogic { retry_partial }, services, health_config, ElasticsearchHealthLogic, 1)(src/sinks/elasticsearch/config.rs 的build)。
部分失败(Partial Failures)
Elasticsearch 默认允许 bulk 部分失败,通常源于索引映射错误(数据键类型不一致)。要改变这一行为,需参考 Elasticsearch 的ignore_malformed设置。默认 Vector 不重试部分失败;开启request_retry_partial后会重试整个部分失败的请求,因此强烈建议配合id_key避免产生重复文档。
Query params 结构
query参数值可以是“单键单值”,也可以是“单键多值”:
sources: source0: query: field: value fruit: - mango - papaya - kiwi源码实现与之一致:QueryParameterValue区分SingleParam与MultiParams两种形态(src/sinks/elasticsearch/common.rs 中 bulk URI 组装时对两者分别以“追加一对”和“同 key 追加多次”处理)。
OpenSearch 兼容性
该 sink 完全兼容 OpenSearch(Elasticsearch 的开源分支),可对接:
- 自管 OpenSearch:与 Elasticsearch 相同配置即可(
opensearch_service_type = "managed",默认值); - Amazon OpenSearch Service:配置 AWS 认证并保持
opensearch_service_type = "managed"; - Amazon OpenSearch Serverless:设置
opensearch_service_type = "serverless"并使用 AWS 认证。
对 OpenSearch 而言,Bulk 索引、Data Streams、Basic/AWS 认证、TLS、API 版本自动检测、压缩与自定义请求头等能力全部可用。
AWS 认证签名
aws-corefeature 启用时,每个请求经sign_request使用 AWS v4 签名(Auth::Aws { credentials_provider, region }),Serverless 模式追加内容哈希头;凭据加载支持 IMDS 与assume_role(load_timeout_secs控制超时)。
健康检查与版本自动检测
- 健康检查:
ElasticsearchCommon::healthcheck向/_cluster/health发起 GET,收到 200 视为健康,其余状态码报UnexpectedStatus;OpenSearch Serverless 不支持健康检查,源码会打印 “Amazon OpenSearch Serverless does not support healthchecks. Skipping healthcheck...” 并直接跳过。 - API 版本自动检测:
get_version请求根路径/,解析响应 JSON 中version.number的主版本号(如7.17.1→ 7)。探测失败时仅记录警告并按启发式假设版本:suppress_type_name = true时假设为 V6,否则假设 V8——源码注释明确说明该假设“远非完美”,并提示未来版本会在配置解析期直接报错。 - 版本对协议的影响:
version >= 7时自动抑制type字段发送(等效旧版suppress_type_name行为)。
端到端确认(Acknowledgements)
acknowledgements控制该 sink 的确认方式(AcknowledgementsConfig,可写布尔或结构体),决定上游 buffer/重试如何感知投递成功。注意组件投递语义为 at-least-once:结合request_retry_partial的重试整包行为,生产环境若开启部分失败重试,务必配合id_key(或 data stream 的create语义 + 稳定type-dataset-namespace)保证幂等性。
测试与验证入口
- 单元/配置校验测试:src/sinks/elasticsearch/tests.rs 与 src/sinks/elasticsearch/config.rs 内嵌的
validate_*系列测试,覆盖端点互斥、非法 URI、非法批量参数、外部版本控制缺参等路径; - 集成测试:src/sinks/elasticsearch/integration_tests.rs(需要
es-integration-testsfeature); - 文档生成测试:
generate_config测试保证ElasticsearchConfig与 website/cue/reference/components/sinks/generated/elasticsearch.cue 的参数文档保持同步。
小结
Vector 的 Elasticsearch sink 以endpoints(多端点 + 故障转移)为骨架,以mode(bulk / data_stream)为双轨:bulk 模式依赖bulk.index模板与index/create/updateaction 控制文档冲突行为;data_stream 模式强制createaction、自动同步data_stream.*字段并重命名时间戳以对齐 ECS。配合api_version自动检测、Basic/AWS 双认证、id_key+ 外部版本控制、request_retry_partial幂等重试,以及默认启用的路由模板 confinement 安全机制,该组件可覆盖自管集群、Elastic 云、Amazon OpenSearch(含 Serverless)三类典型部署形态。所有参数默认值与校验行为均可在 src/sinks/elasticsearch/config.rs、src/sinks/elasticsearch/common.rs 与 website/cue/reference/components/sinks/elasticsearch.cue 中直接查证。
【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考