Skip to content

Instantly share code, notes, and snippets.

@rajanim
Last active June 5, 2019 17:57
Show Gist options
  • Select an option

  • Save rajanim/b614535c7ac6af8973e4e1dddf2494c7 to your computer and use it in GitHub Desktop.

Select an option

Save rajanim/b614535c7ac6af8973e4e1dddf2494c7 to your computer and use it in GitHub Desktop.
import com.google.common.util.concurrent.MoreExecutors;
import com.lucidworks.apollo.component.ExecutorComponent;
import com.lucidworks.apollo.pipeline.*;
import com.lucidworks.apollo.pipeline.async.AbstractAsyncStageConfig;
import com.lucidworks.apollo.pipeline.async.AsyncStage;
import com.lucidworks.apollo.pipeline.async.AsyncStageConfig;
import java.util.concurrent.ExecutorService;
/**
* Convenience base class for stages which need asynchronous processing.
*
* @param <C> The config class, subclass of {@link AbstractAsyncStageConfig}
* @param <R> The intermediate result type, passed between {@link #doProcessing} and {@link #mergeResults}.
* @see AsyncStage
*/
public abstract class AbstractAsyncQueryStage<C extends AbstractAsyncStageConfig, R> extends QueryStage<C>
implements AsyncStage<R, QueryRequestAndResponse> {
private final ExecutorService executorService;
protected AbstractAsyncQueryStage(StageAssistFactoryParams params, ExecutorComponent executorComponent) {
super(params);
if(getConfiguration().isAsyncEnabled()) {
int threadPoolSize = Math.max(0, params.getConfigurationComponent().getInt(
"com.lucidworks.apollo.pipeline.query.async.threadPoolSize", 1000));
if(threadPoolSize == 0) {
executorService = executorComponent.getOrCreateCachedThreadPool("system-async-query-pipelines");
} else {
executorService = executorComponent.getOrCreateFixedThreadPool(threadPoolSize, "system-async-query-pipelines");
}
} else {
executorService = MoreExecutors.newDirectExecutorService();
}
}
@Override
public final QueryRequestAndResponse process(QueryRequestAndResponse message, Context context) throws Exception {
if(getConfiguration().isAsyncEnabled()) {
AsyncStageConfig asyncConfig = getConfiguration().getAsyncConfig();
AsyncStage.runInExecutor(this, message, context, asyncConfig.getAsyncId(), executorService);
} else {
AsyncStage.runInForeground(this, message, context);
}
return message;
}
}
package com.lucidworks.apollo.pipeline.query.stages;
import com.google.common.base.Throwables;
import com.google.inject.assistedinject.Assisted;
import com.lucidworks.apollo.component.ExecutorComponent;
import com.lucidworks.apollo.pipeline.*;
import com.lucidworks.apollo.pipeline.async.AsyncStageRuntimeException;
import com.lucidworks.apollo.pipeline.query.AbstractAsyncQueryStage;
import com.lucidworks.apollo.pipeline.query.QueryRequestAndResponse;
import com.lucidworks.apollo.pipeline.query.config.JDBCQueryStageConfig;
import com.lucidworks.apollo.sql.DataSourceDescriptor;
import com.lucidworks.apollo.blob.sql.SQLHelper;
import com.lucidworks.apollo.pipeline.Context;
import com.lucidworks.apollo.pipeline.StageAssistFactoryParams;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import javax.inject.Inject;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.concurrent.*;
/**
* Connect to a DB and execute a prepared statement and add the results into the request or the context
*/
@AutoDiscover(type = JDBCQueryStageConfig.TYPE)
public class JDBCQueryStage extends AbstractAsyncQueryStage<JDBCQueryStageConfig, ResultSet> {
private transient static Logger LOG = LoggerFactory.getLogger(JDBCQueryStage.class);
private final SQLHelper sqlHelper;
@Inject
public JDBCQueryStage(
@Assisted StageAssistFactoryParams params,
SQLHelper helper,
ExecutorComponent executorComponent) {
super(params, executorComponent);
this.sqlHelper = helper;
}
@Override
public ResultSet doProcessing(QueryRequestAndResponse messageCopy, Context contextCopy) {
JDBCQueryStageConfig config = getConfiguration();
try {
Connection conn = sqlHelper.getConnection(...);
PreparedStatement statement = conn.prepareStatement(config.getPreparedStatement(), ResultSet.TYPE_FORWARD_ONLY,
ResultSet.CONCUR_READ_ONLY);
sqlHelper.bindStatement(statement, config.getPreparedStatementKeys(), contextCopy, messageCopy);
sqlHelper.setRows(statement, config.getRows());
return statement.executeQuery();
} catch (SQLException | ExecutionException e) {
throw new AsyncStageRuntimeException(e);
}
}
@Override
public void mergeResults(ResultSet asyncResult, QueryRequestAndResponse message, Context context) throws Exception {
JDBCQueryStageConfig config = getConfiguration();
try {
String pre = config.getPrefix();
long rowsRetr;
if (config.isJoin()) {
if (config.isAsyncEnabled()) {
LOG.warn("Should not use asynchronous execution in conjunction with 'join' property! Doing it anyways.");
}
rowsRetr = sqlHelper.addResultsToRequest(asyncResult, pre, message.request);
} else {
rowsRetr = sqlHelper.addResultsToContext(asyncResult, pre, context);
}
context.setProperty(SQLHelper.ROWS_RETRIEVED_KEY, rowsRetr);
} finally {
if(asyncResult.getStatement().getConnection() != null) {
asyncResult.getStatement().getConnection().close();
}
if(asyncResult.getStatement() != null) {
asyncResult.getStatement().close();
}
asyncResult.close();
}
}
@Override
public void handleError(Throwable throwable, QueryRequestAndResponse pipelineMessage, Context pipelineContext) throws Exception {
LOG.warn("Error while processing SQL query", throwable);
Throwables.propagateIfPossible(throwable, SQLException.class);
}
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment