Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,131 @@
/*
* Copyright OpenSearch Contributors
* SPDX-License-Identifier: Apache-2.0
*
* The OpenSearch Contributors require contributions made to
* this file be licensed under the Apache-2.0 license or a
* compatible open source license.
*
*/

package org.opensearch.dataprepper.plugins.sink.opensearch;

import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.commons.lang3.StringUtils;
import org.opensearch.client.opensearch._types.VersionType;
import org.opensearch.client.opensearch.core.bulk.BulkOperation;
import org.opensearch.client.opensearch.core.bulk.CreateOperation;
import org.opensearch.client.opensearch.core.bulk.DeleteOperation;
import org.opensearch.client.opensearch.core.bulk.IndexOperation;
import org.opensearch.client.opensearch.core.bulk.UpdateOperation;
import org.opensearch.dataprepper.model.opensearch.OpenSearchBulkActions;
import org.opensearch.dataprepper.plugins.sink.opensearch.bulk.SerializedJson;

import java.io.IOException;
import java.util.Optional;

public class BulkOperationFactory {

private final VersionType versionType;
private final ScriptManager scriptManager;
private final ObjectMapper objectMapper;
private final boolean usingDocumentFilters;

public BulkOperationFactory(final VersionType versionType,
final ScriptManager scriptManager,
final ObjectMapper objectMapper,
final boolean usingDocumentFilters) {
this.versionType = versionType;
this.scriptManager = scriptManager;
this.objectMapper = objectMapper;
this.usingDocumentFilters = usingDocumentFilters;
}

public BulkOperation create(final String action,
Comment thread
dinujoh marked this conversation as resolved.
final SerializedJson document,
final Long version,
final String indexName,
final JsonNode jsonNode) {
final Optional<String> docId = document.getDocumentId();
final Optional<String> routing = document.getRoutingField();
final Optional<String> pipeline = document.getPipelineField();

if (StringUtils.equals(action, OpenSearchBulkActions.CREATE.toString())) {
final CreateOperation.Builder<Object> builder =
new CreateOperation.Builder<>()
.index(indexName)
.document(document);
docId.ifPresent(builder::id);
routing.ifPresent(builder::routing);
pipeline.ifPresent(builder::pipeline);
return new BulkOperation.Builder().create(builder.build()).build();
}

if (StringUtils.equals(action, OpenSearchBulkActions.UPDATE.toString()) ||
StringUtils.equals(action, OpenSearchBulkActions.UPSERT.toString())) {
return createUpdateOperation(action, document, version, indexName, jsonNode, docId, routing);
}

if (StringUtils.equals(action, OpenSearchBulkActions.DELETE.toString())) {
final DeleteOperation.Builder builder = new DeleteOperation.Builder()
.index(indexName)
.versionType(versionType)
.version(version);
docId.ifPresent(builder::id);
routing.ifPresent(builder::routing);
return new BulkOperation.Builder().delete(builder.build()).build();
}

// Default to "index"
final IndexOperation.Builder<Object> builder = new IndexOperation.Builder<>()
.index(indexName)
.document(document)
.version(version)
.versionType(versionType);
docId.ifPresent(builder::id);
routing.ifPresent(builder::routing);
pipeline.ifPresent(builder::pipeline);
return new BulkOperation.Builder().index(builder.build()).build();
}

private BulkOperation createUpdateOperation(final String action,
final SerializedJson document,
final Long version,
final String indexName,
final JsonNode jsonNode,
final Optional<String> docId,
final Optional<String> routing) {
JsonNode filteredJsonNode = jsonNode;
try {
if (usingDocumentFilters) {
filteredJsonNode = objectMapper.reader().readTree(document.getSerializedJson());
}
} catch (final IOException e) {
throw new RuntimeException(
String.format("An exception occurred while deserializing a document for the %s action: %s", action, e.getMessage()));
}

final boolean isUpsert = StringUtils.equals(action.toLowerCase(), OpenSearchBulkActions.UPSERT.toString());
final UpdateOperation.Builder<Object> builder = new UpdateOperation.Builder<>()
.index(indexName)
.versionType(versionType)
.version(version);

if (scriptManager.isScriptEnabled()) {
builder.script(scriptManager.buildScript(filteredJsonNode, document.getResolvedScriptParameters().orElse(null)));
if (isUpsert) {
builder.upsert(filteredJsonNode);
builder.scriptedUpsert(true);
}
} else if (isUpsert) {
builder.document(filteredJsonNode).upsert(filteredJsonNode);
} else {
builder.document(filteredJsonNode);
}

docId.ifPresent(builder::id);
routing.ifPresent(builder::routing);
return new BulkOperation.Builder().update(builder.build()).build();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,6 @@
package org.opensearch.dataprepper.plugins.sink.opensearch;

import com.google.common.annotations.VisibleForTesting;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import io.micrometer.core.instrument.Counter;
import io.micrometer.core.instrument.DistributionSummary;
Expand All @@ -21,10 +20,6 @@
import org.opensearch.client.opensearch._types.VersionType;
import org.opensearch.client.opensearch.core.BulkRequest;
import org.opensearch.client.opensearch.core.bulk.BulkOperation;
import org.opensearch.client.opensearch.core.bulk.CreateOperation;
import org.opensearch.client.opensearch.core.bulk.DeleteOperation;
import org.opensearch.client.opensearch.core.bulk.IndexOperation;
import org.opensearch.client.opensearch.core.bulk.UpdateOperation;
import org.opensearch.client.transport.TransportOptions;
import org.opensearch.common.unit.ByteSizeUnit;
import org.opensearch.dataprepper.aws.api.AwsCredentialsSupplier;
Expand Down Expand Up @@ -140,6 +135,7 @@ public class OpenSearchSink extends AbstractSink<Record<Event>> {
private final String action;
private final List<ActionConfiguration> actions;
private final ScriptManager scriptManager;
private final BulkOperationFactory bulkOperationFactory;
private final String documentRootKey;
private String configuredIndexAlias;
private final ReentrantLock lock;
Expand Down Expand Up @@ -221,6 +217,7 @@ public OpenSearchSink(final PluginSetting pluginSetting,
this.lastFlushTimeMap = new ConcurrentHashMap<>();
this.pluginConfigObservable = pluginConfigObservable;
this.objectMapper = new ObjectMapper();
this.bulkOperationFactory = new BulkOperationFactory(versionType, scriptManager, objectMapper, isUsingDocumentFilters());
this.queryExecutorService = openSearchSinkConfig.getIndexConfiguration().getQueryTerm() != null ?
Executors.newSingleThreadExecutor(BackgroundThreadFactory.defaultExecutorThreadFactory("existing-document-query-manager")) : null;

Expand Down Expand Up @@ -344,96 +341,6 @@ public boolean isReady() {
return initialized;
}

BulkOperation getBulkOperationForAction(final String action,
final SerializedJson document,
final Long version,
final String indexName,
final JsonNode jsonNode) {
BulkOperation bulkOperation;
final Optional<String> docId = document.getDocumentId();
final Optional<String> routing = document.getRoutingField();
final Optional<String> pipeline = document.getPipelineField();

if (StringUtils.equals(action, OpenSearchBulkActions.CREATE.toString())) {
final CreateOperation.Builder<Object> createOperationBuilder =
new CreateOperation.Builder<>()
.index(indexName)
.document(document);
docId.ifPresent(createOperationBuilder::id);
routing.ifPresent(createOperationBuilder::routing);
pipeline.ifPresent(createOperationBuilder::pipeline);

bulkOperation = new BulkOperation.Builder()
.create(createOperationBuilder.build())
.build();
return bulkOperation;
}
if (StringUtils.equals(action, OpenSearchBulkActions.UPDATE.toString()) ||
StringUtils.equals(action, OpenSearchBulkActions.UPSERT.toString())) {

JsonNode filteredJsonNode = jsonNode;
try {
if (isUsingDocumentFilters()) {
filteredJsonNode = objectMapper.reader().readTree(document.getSerializedJson());
}
} catch (final IOException e) {
throw new RuntimeException(
String.format("An exception occurred while deserializing a document for the %s action: %s", action, e.getMessage()));
}


final boolean isUpsert = StringUtils.equals(action.toLowerCase(), OpenSearchBulkActions.UPSERT.toString());
final UpdateOperation.Builder<Object> updateOperationBuilder = new UpdateOperation.Builder<>()
.index(indexName)
.versionType(versionType)
.version(version);

if (scriptManager.isScriptEnabled()) {
updateOperationBuilder.script(scriptManager.buildScript(jsonNode, document.getResolvedScriptParameters().orElse(null)));
updateOperationBuilder.upsert(filteredJsonNode);
updateOperationBuilder.scriptedUpsert(true);
} else if (isUpsert) {
updateOperationBuilder.document(filteredJsonNode).upsert(filteredJsonNode);
} else {
updateOperationBuilder.document(filteredJsonNode);
}

docId.ifPresent(updateOperationBuilder::id);
routing.ifPresent(updateOperationBuilder::routing);
bulkOperation = new BulkOperation.Builder()
.update(updateOperationBuilder.build())
.build();
return bulkOperation;
}
if (StringUtils.equals(action, OpenSearchBulkActions.DELETE.toString())) {
final DeleteOperation.Builder deleteOperationBuilder =
new DeleteOperation.Builder().index(indexName);
docId.ifPresent(deleteOperationBuilder::id);
routing.ifPresent(deleteOperationBuilder::routing);
bulkOperation = new BulkOperation.Builder()
.delete(deleteOperationBuilder
.versionType(versionType)
.version(version)
.build())
.build();
return bulkOperation;
}
// Default to "index"
final IndexOperation.Builder<Object> indexOperationBuilder =
new IndexOperation.Builder<>()
.index(indexName)
.document(document)
.version(version)
.versionType(versionType);
docId.ifPresent(indexOperationBuilder::id);
routing.ifPresent(indexOperationBuilder::routing);
pipeline.ifPresent(indexOperationBuilder::pipeline);
bulkOperation = new BulkOperation.Builder()
.index(indexOperationBuilder.build())
.build();
return bulkOperation;
}

@Override
public void doOutput(final Collection<Record<Event>> records) {
final long threadId = Thread.currentThread().getId();
Expand Down Expand Up @@ -541,7 +448,7 @@ public void doOutput(final Collection<Record<Event>> records) {
BulkOperation bulkOperation;

try {
bulkOperation = getBulkOperationForAction(eventAction, document, version, indexName, event.getJsonNode());
bulkOperation = bulkOperationFactory.create(eventAction, document, version, indexName, event.getJsonNode());
} catch (final Exception e) {
LOG.error("An exception occurred while constructing the bulk operation for a document: ", e);
logFailureForDlqObjects(List.of(createDlqObjectFromEvent(event, indexName, e.getMessage())), e);
Expand Down
Loading
Loading