如何在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
相关产品推荐
相关产品推荐

