You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在AWS OpenSearch摄入管道中将关键词拆分为单独文档?

解决AWS OpenSearch Ingestion Pipeline拆分关键词为单独文档的问题

问题场景

从Parquet文件读取源文档,schema如下:

query
keyword1 keyword2
keyword3 keyword4 keyword5

需求是按空格拆分query字段,将每个关键词作为单独文档写入OpenSearch索引。当前管道配置:

version: "2"
my-pipeline:
  source:
---- removed -----
  processor:
    - split_string:
        entries:
          - source: "query"
            delimiter_regex: "\\s+"
  sink:
    - opensearch:
        hosts: <my host>
        aws: <aws details>
        index: index_1
        document_id_field: query

摄入时出现异常:

Caused by: com.fasterxml.jackson.databind.exc.MismatchedInputException: Cannot deserialize value of type `java.lang.String` from Array value (token `JsonToken.START_ARRAY`)
 at [Source: UNKNOWN; byte offset: #UNKNOWN]
    at com.fasterxml.jackson.databind.exc.MismatchedInputException.from(MismatchedInputException.java:59) ~[jackson-databind-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.databind.DeserializationContext.reportInputMismatch(DeserializationContext.java:1752) ~[jackson-databind-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.databind.DeserializationContext.handleUnexpectedToken(DeserializationContext.java:1526) ~[jackson-databind-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.databind.deser.std.StdDeserializer._deserializeFromArray(StdDeserializer.java:222) ~[jackson-databind-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.databind.deser.std.StringDeserializer.deserialize(StringDeserializer.java:46) ~[jackson-databind-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.databind.deser.std.StringDeserializer.deserialize(StringDeserializer.java:11) ~[jackson-databind-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.databind.deser.DefaultDeserializationContext.readRootValue(DefaultDeserializationContext.java:323) ~[jackson-databind-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.databind.ObjectMapper._readValue(ObjectMapper.java:4801) ~[jackson-databind-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.databind.ObjectMapper.readValue(ObjectMapper.java:2974) ~[jackson-databind-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.databind.ObjectMapper.treeToValue(ObjectMapper.java:3438) ~[jackson-databind-2.15.3.jar:2.15.3]
    at org.opensearch.dataprepper.model.event.JacksonEvent.mapNodeToObject(JacksonEvent.java:209) ~[data-prepper-api-2.6.1.jar:?]
    at org.opensearch.dataprepper.model.event.JacksonEvent.get(JacksonEvent.java:199) ~[data-prepper-api-2.6.1.jar:?]
    at org.opensearch.dataprepper.plugins.sink.opensearch.OpenSearchSink.getDocument(OpenSearchSink.java:468) ~[opensearch-2.6.1.jar:?]
    at org.opensearch.dataprepper.plugins.sink.opensearch.OpenSearchSink.doOutput(OpenSearchSink.java:374) ~[opensearch-2.6.1.jar:?]
    at org.opensearch.dataprepper.model.sink.AbstractSink.lambda$output$0(AbstractSink.java:67) ~[data-prepper-api-2.6.1.jar:?]
    at io.micrometer.core.instrument.composite.CompositeTimer.record(CompositeTimer.java:141) ~[micrometer-core-1.11.5.jar:1.11.5]
    at org.opensearch.dataprepper.model.sink.AbstractSink.output(AbstractSink.java:67) ~[data-prepper-api-2.6.1.jar:?]
    at org.opensearch.dataprepper.pipeline.Pipeline.lambda$publishToSinks$5(Pipeline.java:349) ~[data-prepper-core-2.6.1.jar:?]
    ... 5 more

错误原因

split_string处理器仅将query字段的字符串转换为字符串数组,但未将数组拆分为独立事件。而OpenSearch Sink的document_id_field要求字段值为字符串类型,无法处理数组,因此抛出类型不匹配异常。

解决方案

需要结合split_string和split处理器,先拆分字符串为数组,再将数组拆分为单个事件,最后映射字段确保document_id_field为字符串:

修改后的完整管道配置:

version: "2"
my-pipeline:
  source:
    # 你的Parquet源配置(如S3等)
  processor:
    # 1. 将query字段按空格拆分为字符串数组,存入临时字段keywords
    - split_string:
        entries:
          - source: "query"
            delimiter_regex: "\\s+"
            target: "keywords"
    # 2. 将keywords数组拆分为独立事件,每个事件包含单个关键词字段keyword
    - split:
        source: "keywords"
        target: "keyword"
    # 3. 将单个关键词映射回query字段,保持字段名一致
    - rename:
        entries:
          - from: "keyword"
            to: "query"
    # 4. 清理临时字段keywords(可选)
    - remove_fields:
        fields: ["keywords"]
  sink:
    - opensearch:
        hosts: <my host>
        aws: <aws details>
        index: index_1
        document_id_field: query # 此时query为单个字符串,符合要求

配置说明

  • split_string:将原query的空格分隔字符串转换为数组,存入keywords临时字段,避免覆盖原字段。
  • split:遍历keywords数组,为每个元素生成一个独立事件,每个事件的keyword字段对应单个关键词。
  • rename:将keyword字段重命名为query,确保和原schema字段名一致,也可直接用keyword作为document_id_field。
  • remove_fields:删除临时的keywords字段,清理事件数据(可选步骤)。

配置生效后,每个关键词会作为独立文档写入OpenSearch索引,解决类型不匹配的问题。

内容的提问来源于stack exchange,提问作者ikki

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.01 09:15:32