Revision: 14712
http://gate.svn.sourceforge.net/gate/?rev=14712&view=rev
Author: valyt
Date: 2011-12-09 11:46:30 +0000 (Fri, 09 Dec 2011)
Log Message:
-----------
Local side of RemoteQueryRunner pretty much done.
Modified Paths:
--------------
mimir/trunk/mimir-client/src/gate/mimir/search/RemoteQueryRunner.java
Modified: mimir/trunk/mimir-client/src/gate/mimir/search/RemoteQueryRunner.java
===================================================================
--- mimir/trunk/mimir-client/src/gate/mimir/search/RemoteQueryRunner.java
2011-12-09 11:45:56 UTC (rev 14711)
+++ mimir/trunk/mimir-client/src/gate/mimir/search/RemoteQueryRunner.java
2011-12-09 11:46:30 UTC (rev 14712)
@@ -8,26 +8,27 @@
* Version 3, June 2007 (also included with this distribution as file
* LICENCE-LGPL3.html).
*
- * Valentin Tablan, 05 Jan 2010
+ * Valentin Tablan, 08 Dec 2011
*
* $Id$
*/
package gate.mimir.search;
import gate.mimir.index.IndexException;
+import gate.mimir.index.mg4j.zipcollection.DocumentData;
import gate.mimir.search.query.Binding;
import gate.mimir.tool.WebUtils;
-
import it.unimi.dsi.fastutil.doubles.DoubleArrayList;
-import it.unimi.dsi.fastutil.ints.IntList;
-import it.unimi.dsi.fastutil.objects.ObjectList;
+import it.unimi.dsi.fastutil.ints.Int2ObjectLinkedOpenHashMap;
import java.io.IOException;
import java.io.Serializable;
-import java.util.HashSet;
+import java.net.URLEncoder;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.concurrent.Executor;
import org.apache.log4j.Logger;
@@ -37,28 +38,103 @@
*/
public class RemoteQueryRunner implements QueryRunner {
- protected static final String ACTION_RENDER_DOCUMENT = "renderDocument";
-
protected static final String SERVICE_SEARCH = "search";
- protected static final String ACTION_DOC_TEXT_BIN = "docTextBin";
-
- protected static final String ACTION_DOC_URI_BIN = "docURIBin";
+ protected static final String ACTION_POST_QUERY_BIN = "postQueryBin";
- protected static final String ACTION_DOC_TITLE_BIN = "docTitleBin";
+ protected static final String ACTION_CURRENT_DOC_COUNT_BIN =
"documentsCurrentCountBin";
- protected static final String ACTION_DOC_MEDATADA_FIELDS_BIN =
"docMetadataFieldsBin";
+ protected static final String ACTION_DOC_COUNT_BIN = "documentsCountBin";
+ protected static final String ACTION_DOC_ID_BIN = "documentIdBin";
+
+ protected static final String ACTION_DOC_HITS_BIN = "documentHitsBin";
+
+ protected static final String ACTION_DOC_SCORES_BIN = "documentsScoresBin";
+
+ protected static final String ACTION_DOC_DATA_BIN = "documentDataBin";
+
+ protected static final String ACTION_RENDER_DOCUMENT = "renderDocument";
+
protected static final String ACTION_CLOSE = "close";
+ /**
+ * The maximum number of documents to be stored in the local document cache.
+ */
+ protected static final int DOCUMENT_CACHE_SIZE = 1000;
/**
+ * Action run in a background thread, used to update the document data
+ * (document ID, document score) from the remote endpoint.
+ * This runs once, started during the creation of the query runner.
+ */
+ protected class DocumentDataUpdater implements Runnable {
+ @Override
+ public void run() {
+ // wait for the first pass to complete
+ while(documentsCount < 0) {
+ if(closed) return;
+ // update document counts
+ try {
+ int newDocumentsCount = webUtils.getInt(
+ getActionBaseUrl(ACTION_DOC_COUNT_BIN), "queryId",
+ URLEncoder.encode(queryId, "UTF-8"));
+ if(newDocumentsCount < 0) {
+ // still not finished-> update current count
+ currentDocumentsCount = webUtils.getInt(
+ getActionBaseUrl(ACTION_CURRENT_DOC_COUNT_BIN), "queryId",
+ URLEncoder.encode(queryId, "UTF-8"));
+ // ... and wait a while before asking again
+ Thread.sleep(500);
+ } else {
+ // remote side has finished enumerating all documents
+ // download all the scores in one go.
+ double[] docScores = (double[]) webUtils.getObject(
+ getActionBaseUrl(ACTION_DOC_SCORES_BIN), "queryId",
+ URLEncoder.encode(queryId, "UTF-8"));
+ if(docScores != null && docScores.length > 0) {
+ documentScores = new DoubleArrayList(docScores);
+ if(documentScores.size() != newDocumentsCount) {
+ // malfunction
+ exceptionInBackgroundThread = new RuntimeException(
+ "Incorrect number of document scores from remote side: " +
+ "was expecting " + newDocumentsCount + " but got " +
+ documentScores.size() +"!");
+ return;
+ }
+ } else {
+ // we're not scoring
+ documentScores = null;
+ }
+ // ...and we're done!
+ documentsCount = newDocumentsCount;
+ }
+ } catch(IOException e) {
+ exceptionInBackgroundThread = e;
+ logger.error("Exception while obtaining remote document data", e);
+ } catch(ClassNotFoundException e) {
+ exceptionInBackgroundThread = e;
+ logger.error("Exception while obtaining remote document data", e);
+ } catch(InterruptedException e) {
+ Thread.currentThread().interrupt();
+ logger.warn("Interrupted while waiting", e);
+ }
+ }
+ }
+ }
+
+
+ /**
+ * A cache of MG4J {@link Document}s used for returning the hit text.
+ */
+ private Int2ObjectLinkedOpenHashMap<DocumentData> documentCache;
+
+ /**
* The WebUtils instance we use to communicate with the remote
* index.
*/
private WebUtils webUtils;
-
/**
* The URL to the server hosting the remote index we're searching
*/
@@ -69,8 +145,20 @@
*/
private String queryId;
+ /**
+ * The total number of result documents (or -1 if not yet known).
+ */
+ private volatile int documentsCount;
+ private volatile boolean closed;
+
/**
+ * The current number of documents. After all documents have been retrieved,
+ * this value is identical to {@link #documentsCount}.
+ */
+ private int currentDocumentsCount;
+
+ /**
* Shared Logger
*/
private static Logger logger = Logger.getLogger(RemoteQueryRunner.class);
@@ -83,27 +171,48 @@
*/
private Exception exceptionInBackgroundThread;
-
/**
- * The document IDs for the documents found to contain hits. This list is
- * sorted in ascending documentID order.
- */
- protected IntList documentIds;
-
- /**
* If scoring is enabled ({@link #scorer} is not <code>null</code>), this
list
* contains the scores for the documents found to contain hits. This list is
* aligned to {@link #documentIds}.
*/
protected DoubleArrayList documentScores;
- /**
- * The sets of hits for each returned document. This data structure is
lazily
- * built, so some elements may be null.
- */
- protected ObjectList<List<Binding>> documentHits;
- private String getActionBaseUrl(String action) throws IOException{
+ public RemoteQueryRunner(String remoteUrl, String queryString,
+ Executor threadSource, WebUtils webUtils) throws IOException {
+ this.remoteUrl = remoteUrl.endsWith("/") ? remoteUrl : (remoteUrl + "/");
+ this.webUtils = webUtils;
+ this.closed = false;
+ // submit the remote query
+ try {
+ this.queryId = (String) webUtils.getObject(
+ getActionBaseUrl(ACTION_POST_QUERY_BIN),
+ "queryString",
+ URLEncoder.encode(queryString, "UTF-8"));
+ } catch(ClassNotFoundException e) {
+ //we were expecting a String but got some object of unknown class
+ throw (IOException) new IOException(
+ "Was expecting a String query ID value, but got " +
+ "an unknown object type!").initCause(e);
+ }
+
+ // init the cache
+ documentCache = new Int2ObjectLinkedOpenHashMap<DocumentData>();
+
+ // start the background action
+ documentsCount = -1;
+ currentDocumentsCount = 0;
+ DocumentDataUpdater docDataUpdater = new DocumentDataUpdater();
+ if(threadSource != null) {
+ threadSource.execute(docDataUpdater);
+ } else {
+ new Thread(docDataUpdater,
+ DocumentDataUpdater.class.getCanonicalName()).start();
+ }
+ }
+
+ protected String getActionBaseUrl(String action) throws IOException{
//this method is always called from interactive methods, that are capable
of
//reporting errors to the user. So we use this place to check if the
//background thread had any problems, and report them if so.
@@ -122,22 +231,55 @@
str.append(action);
return str.toString();
}
+
+ /**
+ * Gets (from the cache, or from the remote endpoint) the {@link
DocumentData}
+ * for the document at the specified rank.
+ * @param rank
+ * @return
+ * @throws IndexException
+ * @throws IndexOutOfBoundsException
+ * @throws IOException
+ */
+ protected DocumentData getDocumentData(int rank) throws IndexException,
+ IndexOutOfBoundsException, IOException {
+ DocumentData docData = documentCache.getAndMoveToFirst(rank);
+ if(docData == null) {
+ // cache miss -> remote retrieve
+ try {
+ docData = (DocumentData)webUtils.getObject(
+ getActionBaseUrl(ACTION_DOC_DATA_BIN),
+ "queryId", queryId,
+ "documentRank", Integer.toString(rank));
+ documentCache.putAndMoveToFirst(rank, docData);
+ if(documentCache.size() > DOCUMENT_CACHE_SIZE) {
+ // reduce size
+ documentCache.removeLast();
+ }
+ } catch(IOException e) {
+ throw new IndexException(e);
+ } catch(ClassNotFoundException e) {
+ throw new IndexException("Was expecting a DocumentData value, " +
+ "but got an unknown object type!", e);
+ }
+ }
+ return docData;
+ }
/* (non-Javadoc)
* @see gate.mimir.search.QueryRunner#getDocumentsCount()
*/
@Override
public int getDocumentsCount() {
- // TODO Auto-generated method stub
- return 0;
+ return documentsCount;
}
/* (non-Javadoc)
* @see gate.mimir.search.QueryRunner#getCurrentDocumentsCount()
*/
@Override
- public int getCurrentDocumentsCount() {
- return documentIds.size();
+ public int getDocumentsCurrentCount() {
+ return (documentsCount < 0) ? currentDocumentsCount : documentsCount;
}
/* (non-Javadoc)
@@ -146,7 +288,8 @@
@Override
public int getDocumentID(int rank) throws IndexOutOfBoundsException,
IOException {
- return documentIds.get(rank);
+ return webUtils.getInt(getActionBaseUrl(ACTION_DOC_ID_BIN), "queryId",
+ URLEncoder.encode(queryId, "UTF-8"));
}
/* (non-Javadoc)
@@ -155,16 +298,30 @@
@Override
public double getDocumentScore(int rank) throws IndexOutOfBoundsException,
IOException {
- return documentScores.get(rank);
+ if(documentsCount < 0) {
+ // premature call
+ throw new IndexOutOfBoundsException(
+ "Score requested before collection of documents has completed.");
+ }
+ return documentScores != null ? documentScores.get(rank) : DEFAULT_SCORE;
}
/* (non-Javadoc)
* @see gate.mimir.search.QueryRunner#getDocumentHits(int)
*/
+ @SuppressWarnings("unchecked")
@Override
public List<Binding> getDocumentHits(int rank)
throws IndexOutOfBoundsException, IOException {
- return documentHits.get(rank);
+ try {
+ return (List<Binding>)webUtils.getObject(
+ getActionBaseUrl(ACTION_DOC_HITS_BIN),
+ "queryId", queryId,
+ "documentRank", Integer.toString(rank));
+ } catch(ClassNotFoundException e) {
+ throw new RuntimeException("Got wrong value type from remote endpoint",
+ e);
+ }
}
/* (non-Javadoc)
@@ -173,19 +330,7 @@
@Override
public String[][] getDocumentText(int rank, int termPosition, int length)
throws IndexException, IndexOutOfBoundsException, IOException {
- try {
- return (String[][])webUtils.getObject(
- getActionBaseUrl(ACTION_DOC_TEXT_BIN),
- "queryId", queryId,
- "documentRank", Integer.toString(rank),
- "termPosition", Integer.toString(termPosition),
- "length", Integer.toString(length));
- } catch(IOException e) {
- throw new IndexException(e);
- } catch(ClassNotFoundException e) {
- throw new IndexException("Was expecting a bi-dimensional array of" +
- " Strings, but got an unknown object type!", e);
- }
+ return getDocumentData(rank).getText(termPosition, length);
}
/* (non-Javadoc)
@@ -194,17 +339,7 @@
@Override
public String getDocumentURI(int rank) throws IndexException,
IndexOutOfBoundsException, IOException {
- try {
- return (String)webUtils.getObject(
- getActionBaseUrl(ACTION_DOC_URI_BIN),
- "queryId", queryId,
- "documentRank", Integer.toString(rank));
- } catch(IOException e) {
- throw new IndexException(e);
- } catch(ClassNotFoundException e) {
- throw new IndexException("Was expecting a String value, but got an " +
- "unknown object type!", e);
- }
+ return getDocumentData(rank).getDocumentURI();
}
/* (non-Javadoc)
@@ -213,17 +348,7 @@
@Override
public String getDocumentTitle(int rank) throws IndexException,
IndexOutOfBoundsException, IOException {
- try {
- return (String)webUtils.getObject(
- getActionBaseUrl(ACTION_DOC_TITLE_BIN),
- "queryId", queryId,
- "documentRank", Integer.toString(rank));
- } catch(IOException e) {
- throw new IndexException(e);
- } catch(ClassNotFoundException e) {
- throw new IndexException("Was expecting a String value, but got an " +
- "unknown object type!", e);
- }
+ return getDocumentData(rank).getDocumentTitle();
}
/* (non-Javadoc)
@@ -232,9 +357,7 @@
@Override
public Serializable getDocumentMetadataField(int rank, String fieldName)
throws IndexException, IndexOutOfBoundsException, IOException {
- Set<String> names = new HashSet<String>();
- names.add(fieldName);
- return getDocumentMetadataFields(rank, names).get(fieldName);
+ return getDocumentData(rank).getMetadataField(fieldName);
}
/* (non-Javadoc)
@@ -242,31 +365,15 @@
*/
@Override
public Map<String, Serializable> getDocumentMetadataFields(int rank,
- Set<String> fieldNames) throws IndexException,
- IndexOutOfBoundsException, IOException {
- try {
- // build a comma-separated value
- StringBuilder namesStr = new StringBuilder();
- boolean first = true;
- for(String aName : fieldNames) {
- if(first) {
- first = false;
- } else {
- namesStr.append(", ");
- }
- namesStr.append(aName.replace(",", "\\,"));
- }
- return (Map<String, Serializable>)webUtils.getObject(
- getActionBaseUrl(ACTION_DOC_MEDATADA_FIELDS_BIN),
- "queryId", queryId,
- "documentRank", Integer.toString(rank),
- "fieldNames", namesStr.toString());
- } catch(IOException e) {
- throw new IndexException(e);
- } catch(ClassNotFoundException e) {
- throw new IndexException("Was expecting a Map<String, Serializable> " +
- "value, but got an unknown object type!", e);
- }
+ Set<String> fieldNames) throws IndexException,
IndexOutOfBoundsException,
+ IOException {
+ Map<String, Serializable> res = new HashMap<String, Serializable>();
+ int docId = getDocumentID(rank);
+ for(String fieldName : fieldNames) {
+ Serializable value = getDocumentMetadataField(docId, fieldName);
+ if(value != null) res.put(fieldName, value);
+ }
+ return res;
}
/* (non-Javadoc)
@@ -287,5 +394,7 @@
public void close() throws IOException {
webUtils.getVoid(getActionBaseUrl(ACTION_CLOSE),
"queryId", queryId);
+ closed = true;
+ documentCache.clear();
}
}
This was sent by the SourceForge.net collaborative development platform, the
world's largest Open Source development site.
------------------------------------------------------------------------------
Cloud Services Checklist: Pricing and Packaging Optimization
This white paper is intended to serve as a reference, checklist and point of
discussion for anyone considering optimizing the pricing and packaging model
of a cloud services business. Read Now!
http://www.accelacomm.com/jaw/sfnl/114/51491232/
_______________________________________________
GATE-cvs mailing list
[email protected]
https://lists.sourceforge.net/lists/listinfo/gate-cvs