Last active
June 5, 2019 17:57
-
-
Save rajanim/b614535c7ac6af8973e4e1dddf2494c7 to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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; | |
| } | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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