要解决Apache NiFi QuerySolr处理器达到"start"参数上限的问题,可以使用分页查询来处理大量数据。下面是一个示例代码:
import org.apache.nifi.processor.AbstractProcessor;
import org.apache.nifi.processor.ProcessContext;
import org.apache.nifi.processor.ProcessSession;
import org.apache.nifi.processor.ProcessSessionFactory;
import org.apache.nifi.processor.Relationship;
import org.apache.nifi.processor.exception.ProcessException;
import org.apache.nifi.processor.io.InputStreamCallback;
import org.apache.nifi.processor.io.OutputStreamCallback;
import org.apache.nifi.processor.util.StandardValidators;
import org.apache.solr.client.solrj.SolrClient;
import org.apache.solr.client.solrj.SolrQuery;
import org.apache.solr.client.solrj.SolrQuery.ORDER;
import org.apache.solr.client.solrj.SolrResponse;
import org.apache.solr.client.solrj.SolrServerException;
import org.apache.solr.client.solrj.impl.HttpSolrClient;
import org.apache.solr.client.solrj.response.QueryResponse;
import org.apache.solr.common.SolrDocument;
import org.apache.solr.common.SolrDocumentList;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
public class QuerySolrProcessor extends AbstractProcessor {
public static final Relationship SUCCESS = new Relationship.Builder()
.name("success")
.description("Success relationship")
.build();
public static final Relationship FAILURE = new Relationship.Builder()
.name("failure")
.description("Failure relationship")
.build();
public static final PropertyDescriptor SOLR_URL = new PropertyDescriptor
.Builder().name("Solr URL")
.description("URL of the Solr server")
.required(true)
.addValidator(StandardValidators.URL_VALIDATOR)
.build();
public static final PropertyDescriptor QUERY = new PropertyDescriptor
.Builder().name("Query")
.description("Solr query string")
.required(true)
.addValidator(StandardValidators.NON_EMPTY_VALIDATOR)
.build();
public static final PropertyDescriptor PAGE_SIZE = new PropertyDescriptor
.Builder().name("Page Size")
.description("Number of documents to retrieve per page")
.required(true)
.defaultValue("1000")
.addValidator(StandardValidators.POSITIVE_INTEGER_VALIDATOR)
.build();
public static final PropertyDescriptor MAX_RESULTS = new PropertyDescriptor
.Builder().name("Max Results")
.description("Maximum number of documents to retrieve")
.required(true)
.defaultValue("10000")
.addValidator(StandardValidators.POSITIVE_INTEGER_VALIDATOR)
.build();
public static final PropertyDescriptor OUTPUT_FIELD = new PropertyDescriptor
.Builder().name("Output Field")
.description("Name of the field in the Solr documents to output")
.required(true)
.addValidator(StandardValidators.NON_EMPTY_VALIDATOR)
.build();
private List descriptors;
private Set relationships;
public QuerySolrProcessor() {
List descriptors = new ArrayList<>();
descriptors.add(SOLR_URL);
descriptors.add(QUERY);
descriptors.add(PAGE_SIZE);
descriptors.add(MAX_RESULTS);
descriptors.add(OUTPUT_FIELD);
this.descriptors = Collections.unmodifiableList(descriptors);
Set relationships = new HashSet<>();
relationships.add(SUCCESS);
relationships.add(FAILURE);
this.relationships = Collections.unmodifiableSet(relationships);
}
@Override
public Set getRelationships() {
return relationships;
}
@Override
public final List getSupportedPropertyDescriptors() {
return descriptors;
}
@Override
public void onTrigger(ProcessContext context, ProcessSessionFactory sessionFactory) throws ProcessException {
final ProcessSession session = sessionFactory.createSession();
final String solrUrl = context.getProperty(SOLR_URL).getValue();
final String query = context.getProperty(QUERY).getValue();
final int pageSize = context.getProperty(PAGE_SIZE).asInteger();
final int maxResults = context.getProperty(MAX_RESULTS).asInteger();
final String outputField = context.getProperty(OUTPUT_FIELD).getValue();
try {
SolrClient solrClient = new HttpSolrClient.Builder(solrUrl).build();
int start = 0;
int retrievedResults = 0;
while (retrievedResults < maxResults) {
SolrQuery solrQuery = new SolrQuery(query);
solrQuery.setStart(start);
solrQuery.setRows(pageSize);
solrQuery.setSort(outputField