diff --git a/contrib/src/main/java/org/archive/modules/postprocessor/TroughCrawlLogFeed.java b/contrib/src/main/java/org/archive/modules/postprocessor/TroughCrawlLogFeed.java index 492a9aa7..a2e3e1bf 100644 --- a/contrib/src/main/java/org/archive/modules/postprocessor/TroughCrawlLogFeed.java +++ b/contrib/src/main/java/org/archive/modules/postprocessor/TroughCrawlLogFeed.java @@ -18,27 +18,21 @@ */ package org.archive.modules.postprocessor; -import java.io.IOException; -import java.nio.charset.Charset; +import java.net.MalformedURLException; import java.util.ArrayList; import java.util.Date; import java.util.List; +import java.util.logging.Level; import java.util.logging.Logger; import org.apache.commons.collections.Closure; -import org.apache.commons.io.IOUtils; -import org.apache.commons.lang.StringUtils; -import org.apache.http.HttpResponse; -import org.apache.http.client.methods.HttpPost; -import org.apache.http.entity.StringEntity; -import org.apache.http.impl.client.CloseableHttpClient; -import org.apache.http.impl.client.HttpClientBuilder; import org.archive.crawler.framework.Frontier; import org.archive.crawler.frontier.AbstractFrontier; import org.archive.crawler.frontier.BdbFrontier; import org.archive.modules.CrawlURI; import org.archive.modules.Processor; import org.archive.modules.net.ServerCache; +import org.archive.spring.KeyedProperties; import org.archive.trough.TroughClient; import org.archive.util.MimetypeUtils; import org.springframework.beans.factory.annotation.Autowired; @@ -98,9 +92,42 @@ public class TroughCrawlLogFeed extends Processor implements Lifecycle { protected static final int BATCH_MAX_TIME_MS = 20 * 1000; protected static final int BATCH_MAX_SIZE = 400; - protected List crawledBatch = new ArrayList(); + protected KeyedProperties kp = new KeyedProperties(); + public KeyedProperties getKeyedProperties() { + return kp; + } + + public void setSegmentId(String segmentId) { + kp.put("segmentId", segmentId); + } + public String getSegmentId() { + return (String) kp.get("segmentId"); + } + + /** + * @param rethinkUrl url with schema rethinkdb:// pointing to + * trough configuration database + */ + public void setRethinkUrl(String rethinkUrl) { + kp.put("rethinkUrl", rethinkUrl); + } + public String getRethinkUrl() { + return (String) kp.get("rethinkUrl"); + } + + protected TroughClient troughClient = null; + + protected TroughClient troughClient() throws MalformedURLException { + if (troughClient == null) { + troughClient = new TroughClient(getRethinkUrl(), 60 * 60); + troughClient.start(); + } + return troughClient; + } + + protected List crawledBatch = new ArrayList(); protected long crawledBatchLastTime = System.currentTimeMillis(); - protected List uncrawledBatch = new ArrayList(); + protected List uncrawledBatch = new ArrayList(); protected long uncrawledBatchLastTime = System.currentTimeMillis(); protected Frontier frontier; @@ -122,14 +149,6 @@ public class TroughCrawlLogFeed extends Processor implements Lifecycle { this.serverCache = serverCache; } - protected String writeUrl = null; - public void setWriteUrl(String writeUrl) { - this.writeUrl = writeUrl; - } - public String getWriteUrl() { - return writeUrl; - } - @Override protected boolean shouldProcess(CrawlURI curi) { if (frontier instanceof AbstractFrontier) { @@ -170,31 +189,6 @@ public class TroughCrawlLogFeed extends Processor implements Lifecycle { super.stop(); } - transient protected CloseableHttpClient _httpClient; - protected CloseableHttpClient httpClient() { - if (_httpClient == null) { - _httpClient = HttpClientBuilder.create().build(); - } - return _httpClient; - } - - protected void post(String statement) { - HttpPost httpPost = new HttpPost(writeUrl); - try { - httpPost.setEntity(new StringEntity(statement)); - HttpResponse response = httpClient().execute(httpPost); - String payload = IOUtils.toString(response.getEntity().getContent(), Charset.forName("UTF-8")); - if (response.getStatusLine().getStatusCode() != 200) { - throw new IOException("expected 200 OK but got " + response.getStatusLine() + " - " + payload); - } - } catch (Exception e) { - logger.warning("problem posting batch to " + writeUrl + " - " + e); - } finally { - httpPost.releaseConnection(); - } - - } - @Override protected void innerProcess(CrawlURI curi) throws InterruptedException { if (curi.getFetchStatus() > 0) { @@ -209,69 +203,111 @@ public class TroughCrawlLogFeed extends Processor implements Lifecycle { } else { warcContentBytes = 0; } + + Object[] values = new Object[] { + new Date(curi.getFetchBeginTime()), + curi.getFetchStatus(), + curi.getContentSize(), + curi.getContentLength(), + curi, + curi.getPathFromSeed(), + (curi.isSeed() && !"".equals(curi.getPathFromSeed())) ? 1 : 0, + curi.getVia(), + MimetypeUtils.truncate(curi.getContentType()), + curi.getContentDigestSchemeString(), + curi.getSourceTag(), + curi.isRevisit() ? 1 : 0, + curi.getExtraInfo().opt("warcFilename"), + curi.getExtraInfo().opt("warcOffset"), + warcContentBytes, + serverCache.getHostFor(curi.getUURI()).getHostName(), + }; + synchronized (crawledBatch) { - crawledBatch.add("(" + TroughClient.sqlValue(new Date(curi.getFetchBeginTime())) + ", " - + TroughClient.sqlValue((Object) curi.getFetchStatus()) + ", " - + TroughClient.sqlValue((Object) curi.getContentSize()) + ", " - + TroughClient.sqlValue((Object) curi.getContentLength()) + ", " - + TroughClient.sqlValue(curi) + ", " - + TroughClient.sqlValue(curi.getPathFromSeed()) + ", " - + TroughClient.sqlValue((Object) (curi.isSeed() && !"".equals(curi.getPathFromSeed()) ? 1 : 0)) + ", " - + TroughClient.sqlValue(curi.getVia()) + ", " - + TroughClient.sqlValue(MimetypeUtils.truncate(curi.getContentType())) + ", " - + TroughClient.sqlValue(curi.getContentDigestSchemeString()) + ", " - + TroughClient.sqlValue(curi.getSourceTag()) + ", " - + TroughClient.sqlValue((Object) (curi.isRevisit() ? 1 : 0)) + ", " - + TroughClient.sqlValue(curi.getExtraInfo().opt("warcFilename")) + ", " - + TroughClient.sqlValue(curi.getExtraInfo().opt("warcOffset")) + ", " - + TroughClient.sqlValue((Object) warcContentBytes) + ", " - + TroughClient.sqlValue(serverCache.getHostFor(curi.getUURI()).getHostName()) + ")"); - if (crawledBatch.size() >= BATCH_MAX_SIZE || System.currentTimeMillis() - crawledBatchLastTime > BATCH_MAX_TIME_MS) { - postCrawledBatch(); - } + crawledBatch.add(values); + } + + if (crawledBatch.size() >= BATCH_MAX_SIZE || System.currentTimeMillis() - crawledBatchLastTime > BATCH_MAX_TIME_MS) { + postCrawledBatch(); } } else { + Object[] values = new Object[] { + new Date(), + curi, + curi.getPathFromSeed(), + curi.getFetchStatus(), + curi.getVia(), + curi.getSourceTag(), + serverCache.getHostFor(curi.getUURI()).getHostName(), + }; + synchronized (uncrawledBatch) { - uncrawledBatch.add("(" - + TroughClient.sqlValue(new Date()) + ", " - + TroughClient.sqlValue(curi) + ", " - + TroughClient.sqlValue(curi.getPathFromSeed()) + ", " - + TroughClient.sqlValue((Object) curi.getFetchStatus()) + ", " - + TroughClient.sqlValue(curi.getVia()) + ", " - + TroughClient.sqlValue(curi.getSourceTag()) + ", " - + TroughClient.sqlValue(serverCache.getHostFor(curi.getUURI()).getHostName()) + ")"); - if (uncrawledBatch.size() >= BATCH_MAX_SIZE || System.currentTimeMillis() - uncrawledBatchLastTime > BATCH_MAX_TIME_MS) { - postUncrawledBatch(); - } + uncrawledBatch.add(values); + } + + if (uncrawledBatch.size() >= BATCH_MAX_SIZE || System.currentTimeMillis() - uncrawledBatchLastTime > BATCH_MAX_TIME_MS) { + postUncrawledBatch(); } } } protected void postCrawledBatch() { - logger.info("posting batch of " + crawledBatch.size() + " crawled urls to " + writeUrl); + logger.info("posting batch of " + crawledBatch.size() + " crawled urls trough segment " + getSegmentId()); synchronized (crawledBatch) { - String sql = "insert into crawled_url (" - + "timestamp, status_code, size, payload_size, url, hop_path, is_seed_redirect, " - + "via, mimetype, content_digest, seed, is_duplicate, warc_filename, " - + "warc_offset, warc_content_bytes, host) values " - + StringUtils.join(crawledBatch, ", ") - + ";"; - post(sql); - crawledBatchLastTime = System.currentTimeMillis(); - crawledBatch.clear(); + if (!crawledBatch.isEmpty()) { + StringBuffer sqlTmpl = new StringBuffer(); + sqlTmpl.append("insert into crawled_url (" + + "timestamp, status_code, size, payload_size, url, hop_path, is_seed_redirect, " + + "via, mimetype, content_digest, seed, is_duplicate, warc_filename, " + + "warc_offset, warc_content_bytes, host) values " + + "(%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)"); + for (int i = 1; i < crawledBatch.size(); i++) { + sqlTmpl.append(", (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)"); + } + + Object[] flattenedValues = new Object[16 * crawledBatch.size()]; + for (int i = 0; i < crawledBatch.size(); i++) { + System.arraycopy(crawledBatch.get(i), 0, flattenedValues, 16 * i, 16); + } + + try { + troughClient().write(getSegmentId(), sqlTmpl.toString(), flattenedValues); + } catch (Exception e) { + logger.log(Level.WARNING, "problem posting batch of " + crawledBatch.size() + " crawled urls to trough segment " + getSegmentId(), e); + } + + crawledBatchLastTime = System.currentTimeMillis(); + crawledBatch.clear(); + } } } - protected void postUncrawledBatch() { - logger.info("posting batch of " + uncrawledBatch.size() + " uncrawled urls to " + writeUrl); + logger.info("posting batch of " + uncrawledBatch.size() + " uncrawled urls trough segment " + getSegmentId()); synchronized (uncrawledBatch) { - String sql = "insert into uncrawled_url (" - + "timestamp, url, hop_path, status_code, via, seed, host) values " - + StringUtils.join(uncrawledBatch, ", ") - + ";"; - post(sql); - uncrawledBatchLastTime = System.currentTimeMillis(); - uncrawledBatch.clear(); + if (!uncrawledBatch.isEmpty()) { + StringBuffer sqlTmpl = new StringBuffer(); + sqlTmpl.append( + "insert into uncrawled_url (timestamp, url, hop_path, status_code, via, seed, host)" + + " values (%s, %s, %s, %s, %s, %s, %s)"); + + for (int i = 1; i < uncrawledBatch.size(); i++) { + sqlTmpl.append(", (%s, %s, %s, %s, %s, %s, %s)"); + } + + Object[] flattenedValues = new Object[7 * uncrawledBatch.size()]; + for (int i = 0; i < uncrawledBatch.size(); i++) { + System.arraycopy(uncrawledBatch.get(i), 0, flattenedValues, 7 * i, 7); + } + + try { + troughClient().write(getSegmentId(), sqlTmpl.toString(), flattenedValues); + } catch (Exception e) { + logger.log(Level.WARNING, "problem posting batch of " + uncrawledBatch.size() + " uncrawled urls to trough segment " + getSegmentId(), e); + } + + uncrawledBatchLastTime = System.currentTimeMillis(); + uncrawledBatch.clear(); + } } } }