From ada4ff913e947a33ebc7621873e59c71d224bcf7 Mon Sep 17 00:00:00 2001 From: David Hoyt Date: Wed, 31 May 2017 15:12:20 -0700 Subject: [PATCH 1/3] Adjust .gitignore for IntelliJ. --- .gitignore | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/.gitignore b/.gitignore index 06cd5ba..aa22773 100644 --- a/.gitignore +++ b/.gitignore @@ -1,8 +1,16 @@ +!.gitignore + +out/ build/ +gradle/ .gradle/ +.idea/ *~ *.jar cassandra-loader *.BADINSERT *.BADPARSE *.LOG +*.iml +/gradlew +/gradlew.bat From 254c45b8d0b957622ef7d5577b166821e95f7eb9 Mon Sep 17 00:00:00 2001 From: David Hoyt Date: Wed, 31 May 2017 15:13:17 -0700 Subject: [PATCH 2/3] Have the CSV parser emit null values as the default null string so we can pick up and detect alternative null strings later. --- src/main/java/com/datastax/loader/parser/DelimParser.java | 1 + 1 file changed, 1 insertion(+) diff --git a/src/main/java/com/datastax/loader/parser/DelimParser.java b/src/main/java/com/datastax/loader/parser/DelimParser.java index 5641549..dd8c14d 100644 --- a/src/main/java/com/datastax/loader/parser/DelimParser.java +++ b/src/main/java/com/datastax/loader/parser/DelimParser.java @@ -86,6 +86,7 @@ public DelimParser(String inDelimiter, int inCharsPerColumn, settings.setKeepQuotes(true); settings.setKeepEscapeSequences(true); settings.getFormat().setComment(comment); + settings.setNullValue(DEFAULT_NULLSTRING); csvp = new CsvParser(settings); } From 0647e3017de9e2c4a15241ee1da755ca01734714 Mon Sep 17 00:00:00 2001 From: David Hoyt Date: Wed, 31 May 2017 16:43:47 -0700 Subject: [PATCH 3/3] Support multiple blob formats (base-64 and hex string). --- build.gradle | 2 +- .../com/datastax/loader/CqlDelimLoad.java | 38 ++++++------ .../com/datastax/loader/CqlDelimLoadTask.java | 16 +++-- .../com/datastax/loader/CqlDelimParser.java | 27 ++++---- .../com/datastax/loader/CqlDelimUnload.java | 32 +++++++--- .../loader/parser/ByteBufferParser.java | 61 +++++++++++++++++-- 6 files changed, 124 insertions(+), 52 deletions(-) diff --git a/build.gradle b/build.gradle index 4813e2b..234dc62 100644 --- a/build.gradle +++ b/build.gradle @@ -1,7 +1,7 @@ apply plugin: 'java' apply plugin: 'application' -def versionNum = '0.0.27' +def versionNum = '0.0.28' allprojects { tasks.withType(JavaCompile) { diff --git a/src/main/java/com/datastax/loader/CqlDelimLoad.java b/src/main/java/com/datastax/loader/CqlDelimLoad.java index d7d98fb..cc230c3 100644 --- a/src/main/java/com/datastax/loader/CqlDelimLoad.java +++ b/src/main/java/com/datastax/loader/CqlDelimLoad.java @@ -16,11 +16,7 @@ package com.datastax.loader; import com.datastax.loader.parser.BooleanParser; -import com.datastax.loader.futures.FutureManager; -import com.datastax.loader.futures.PrintingFutureSet; -import java.util.List; -import java.util.ArrayList; import java.util.Map; import java.util.HashMap; import java.util.Set; @@ -35,11 +31,9 @@ import java.util.concurrent.Callable; import java.util.concurrent.Future; import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; import java.io.File; import java.io.BufferedReader; import java.io.FileReader; -import java.io.InputStreamReader; import java.io.IOException; import java.io.FileOutputStream; import java.io.InputStream; @@ -48,10 +42,6 @@ import java.io.PrintStream; import java.io.FileNotFoundException; import java.text.ParseException; -import java.nio.file.Files; -import java.nio.file.Path; -import java.nio.file.Paths; -import java.nio.file.StandardCopyOption; import java.security.KeyStore; import java.security.SecureRandom; import java.security.KeyStoreException; @@ -68,22 +58,17 @@ import com.datastax.driver.core.Session; import com.datastax.driver.core.ProtocolVersion; import com.datastax.driver.core.PoolingOptions; -import com.datastax.driver.core.ProtocolOptions; import com.datastax.driver.core.HostDistance; -import com.datastax.driver.core.PreparedStatement; import com.datastax.driver.core.ConsistencyLevel; -import com.datastax.driver.core.BoundStatement; -import com.datastax.driver.core.ResultSetFuture; import com.datastax.driver.core.SSLOptions; import com.datastax.driver.core.RemoteEndpointAwareJdkSSLOptions; import com.datastax.driver.core.policies.TokenAwarePolicy; import com.datastax.driver.core.policies.DCAwareRoundRobinPolicy; -import com.codahale.metrics.Timer; -import org.apache.commons.lang3.StringEscapeUtils; +import com.datastax.loader.parser.ByteBufferParser; public class CqlDelimLoad { - private String version = "0.0.27"; + private String version = "0.0.28"; private String host = null; private int port = 9042; private String username = null; @@ -127,6 +112,7 @@ public class CqlDelimLoad { private Locale locale = null; private BooleanParser.BoolStyle boolStyle = null; + private ByteBufferParser.BlobFormat blobFormat = null; private String dateFormatString = null; private String localDateFormatString = "yyyy-MM-dd"; private String nullString = null; @@ -169,6 +155,7 @@ private String usage() { usage.append(" -batchSize Number of INSERTs to batch together [1]\n"); usage.append(" -decimalDelim Decimal delimiter [.] Other option is ','\n"); usage.append(" -boolStyle Style for booleans [TRUE_FALSE]\n"); + usage.append(" -blobFormat [base64|hex] Format of blobs: base-64 encoded or hex string [base64]\n"); usage.append(" -numThreads Number of concurrent threads (files) to load [num cores]\n"); usage.append(" -queryTimeout <# seconds> Query timeout (in seconds) [2]\n"); usage.append(" -numRetries Number of times to retry the INSERT [1]\n"); @@ -178,7 +165,7 @@ private String usage() { usage.append(" -rateFile Where to print the rate statistics\n"); usage.append(" -successDir Directory where to move successfully loaded files\n"); usage.append(" -failureDir Directory where to move files that did not successfully load\n"); - usage.append(" -nullsUnset [false|true] Treat nulls as unset [faslse]\n"); + usage.append(" -nullsUnset [false|true] Treat nulls as unset [false]\n"); usage.append(" -format [delim|jsonline|jsonarray] Format of data: delimited or JSON [delim]\n"); usage.append(" -table Table name (when using JSON)\n"); usage.append(" -keyspace Keyspace name (when using JSON)\n"); @@ -448,6 +435,13 @@ private boolean parseArgs(String[] args) throws IOException, FileNotFoundExcepti return false; } } + if (null != (tkey = amap.remove("-blobFormat"))) { + blobFormat = ByteBufferParser.getBlobFormat(tkey); + if (null == blobFormat) { + System.err.println("Bad blob format. Options are: " + ByteBufferParser.getOptions()); + return false; + } + } if (null != (tkey = amap.remove("-nullsUnset"))) nullsUnset = Boolean.parseBoolean(tkey); if (null != (tkey = amap.remove("-charsPerColumn"))) charsPerColumn = Integer.parseInt(tkey); @@ -644,7 +638,9 @@ public int compare(File f1, File f2) { commentString, dateFormatString, localDateFormatString, - boolStyle, locale, + boolStyle, + blobFormat, + locale, maxErrors, skipRows, skipCols, maxRows, badDir, infile, @@ -669,7 +665,9 @@ public int compare(File f1, File f2) { commentString, dateFormatString, localDateFormatString, - boolStyle, locale, + boolStyle, + blobFormat, + locale, maxErrors, skipRows, skipCols, maxRows, badDir, tFile, diff --git a/src/main/java/com/datastax/loader/CqlDelimLoadTask.java b/src/main/java/com/datastax/loader/CqlDelimLoadTask.java index 1e45446..6d093c7 100644 --- a/src/main/java/com/datastax/loader/CqlDelimLoadTask.java +++ b/src/main/java/com/datastax/loader/CqlDelimLoadTask.java @@ -25,6 +25,7 @@ import com.datastax.loader.futures.PrintingFutureSet; import com.datastax.loader.futures.JsonPrintingFutureSet; import com.datastax.loader.parser.BooleanParser; +import com.datastax.loader.parser.ByteBufferParser; import org.json.simple.JSONArray; import org.json.simple.JSONObject; import org.json.simple.parser.JSONParser; @@ -84,6 +85,7 @@ class CqlDelimLoadTask implements Callable { private String cqlSchema; private Locale locale = null; private BooleanParser.BoolStyle boolStyle = null; + private ByteBufferParser.BlobFormat blobFormat = null; private String dateFormatString = null; private String localDateFormatString = null; private String nullString = null; @@ -105,7 +107,8 @@ public CqlDelimLoadTask(String inCqlSchema, String inDelimiter, int inCharsPerColumn, String inNullString, String inCommentString, String inDateFormatString, String inLocalDateFormatString, - BooleanParser.BoolStyle inBoolStyle, + BooleanParser.BoolStyle inBoolStyle, + ByteBufferParser.BlobFormat inBlobFormat, Locale inLocale, long inMaxErrors, long inSkipRows, String inSkipCols, long inMaxRows, @@ -125,6 +128,7 @@ public CqlDelimLoadTask(String inCqlSchema, String inDelimiter, dateFormatString = inDateFormatString; localDateFormatString = inLocalDateFormatString; boolStyle = inBoolStyle; + blobFormat = inBlobFormat; locale = inLocale; maxErrors = inMaxErrors; skipRows = inSkipRows; @@ -189,16 +193,18 @@ private void setup() throws IOException, ParseException, org.json.simple.parser. cdp = new CqlDelimParser(cqlSchema, delimiter, charsPerColumn, nullString, commentString, dateFormatString, localDateFormatString, - boolStyle, locale, - skipCols, session, true, ttl); + boolStyle, blobFormat, + locale, skipCols, + session, true, ttl); } else if (format.equalsIgnoreCase("jsonline") || format.equalsIgnoreCase("jsonarray")) { cdp = new CqlDelimParser(keyspace, table, delimiter, charsPerColumn, nullString, commentString, dateFormatString, localDateFormatString, - boolStyle, locale, - skipCols, session, true, ttl); + boolStyle, blobFormat, + locale, skipCols, + session, true, ttl); } insert = cdp.generateInsert(); diff --git a/src/main/java/com/datastax/loader/CqlDelimParser.java b/src/main/java/com/datastax/loader/CqlDelimParser.java index c570574..7c9bb82 100644 --- a/src/main/java/com/datastax/loader/CqlDelimParser.java +++ b/src/main/java/com/datastax/loader/CqlDelimParser.java @@ -50,8 +50,6 @@ import java.util.Locale; import java.util.Map; import java.util.Set; -import java.util.regex.Matcher; -import java.util.regex.Pattern; import org.json.simple.JSONObject; import org.json.simple.parser.JSONParser; @@ -72,13 +70,15 @@ public class CqlDelimParser { public CqlDelimParser(String inCqlSchema, String inDelimiter, int inCharsPerColumn, String inNullString, String inCommentString, String inDateFormatString, String inLocalDateFormatString, - BooleanParser.BoolStyle inBoolStyle, Locale inLocale, - String skipList, Session session, boolean bLoader, int inTtl) + BooleanParser.BoolStyle inBoolStyle, ByteBufferParser.BlobFormat inBlobFormat, + Locale inLocale, String skipList, + Session session, boolean bLoader, int inTtl) throws ParseException { // Optionally provide things for the line parser - date format, boolean format, locale ttl = inTtl; - initPmap(inDateFormatString, inLocalDateFormatString, inBoolStyle, - inLocale, bLoader); + initPmap(inDateFormatString, inLocalDateFormatString, + inBoolStyle, inBlobFormat, + inLocale, bLoader); processCqlSchema(inCqlSchema, session); createDelimParser(inDelimiter, inCharsPerColumn, inNullString, inCommentString, skipList); } @@ -87,15 +87,17 @@ public CqlDelimParser(String inKeyspace, String inTable, String inDelimiter, int inCharsPerColumn, String inNullString, String inCommentString, String inDateFormatString, String inLocalDateFormatString, - BooleanParser.BoolStyle inBoolStyle, Locale inLocale, - String skipList, Session session, boolean bLoader, int inTtl) + BooleanParser.BoolStyle inBoolStyle, ByteBufferParser.BlobFormat inBlobFormat, + Locale inLocale, String skipList, + Session session, boolean bLoader, int inTtl) throws ParseException { // Optionally provide things for the line parser - date format, boolean format, locale ttl = inTtl; keyspace = inKeyspace; tablename = inTable; - initPmap(inDateFormatString, inLocalDateFormatString, inBoolStyle, - inLocale, bLoader); + initPmap(inDateFormatString, inLocalDateFormatString, + inBoolStyle, inBlobFormat, + inLocale, bLoader); processCqlSchema(session); createDelimParser(inDelimiter, inCharsPerColumn, inNullString, inCommentString, skipList); } @@ -117,7 +119,8 @@ private class SchemaBits { // intialize the Parsers and the parser map private void initPmap(String dateFormatString, String localDateFormatString, - BooleanParser.BoolStyle inBoolStyle, + BooleanParser.BoolStyle inBoolStyle, + ByteBufferParser.BlobFormat inBlobFormat, Locale inLocale, boolean bLoader) { pmap = new HashMap(); Parser byteParser = new ByteParser(inLocale, bLoader); @@ -131,7 +134,7 @@ private void initPmap(String dateFormatString, String localDateFormatString, Parser uuidParser = new UUIDParser(); Parser bigDecimalParser = new BigDecimalParser(); Parser bigIntegerParser = new BigIntegerParser(); - Parser byteBufferParser = new ByteBufferParser(); + Parser byteBufferParser = new ByteBufferParser(inBlobFormat); Parser inetAddressParser = new InetAddressParser(); Parser dateParser = new DateParser(dateFormatString); Parser localDateParser = new LocalDateParser(localDateFormatString); diff --git a/src/main/java/com/datastax/loader/CqlDelimUnload.java b/src/main/java/com/datastax/loader/CqlDelimUnload.java index 472e33b..7ee03a3 100644 --- a/src/main/java/com/datastax/loader/CqlDelimUnload.java +++ b/src/main/java/com/datastax/loader/CqlDelimUnload.java @@ -23,15 +23,12 @@ import java.util.HashMap; import java.util.Set; import java.util.HashSet; -import java.util.Deque; -import java.util.ArrayDeque; import java.util.Locale; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Callable; import java.util.concurrent.Future; import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; import java.math.BigInteger; import java.io.FileOutputStream; import java.io.BufferedOutputStream; @@ -39,7 +36,6 @@ import java.io.File; import java.io.BufferedReader; import java.io.FileReader; -import java.io.InputStreamReader; import java.io.InputStream; import java.io.FileInputStream; import java.io.IOException; @@ -71,10 +67,11 @@ import com.datastax.driver.core.policies.TokenAwarePolicy; import com.datastax.driver.core.policies.DCAwareRoundRobinPolicy; import com.datastax.driver.core.exceptions.QueryValidationException; +import com.datastax.loader.parser.ByteBufferParser; public class CqlDelimUnload { - private String version = "0.0.27"; + private String version = "0.0.28"; private String host = null; private int port = 9042; private String username = null; @@ -96,6 +93,7 @@ public class CqlDelimUnload { private Locale locale = null; private BooleanParser.BoolStyle boolStyle = null; + private ByteBufferParser.BlobFormat blobFormat = null; private String dateFormatString = null; private String localDateFormatString = "yyyy-MM-dd"; private String nullString = null; @@ -124,6 +122,7 @@ private String usage() { usage.append(" -consistencyLevel Consistency level [LOCAL_ONE]\n"); usage.append(" -decimalDelim Decimal delimiter [.] Other option is ','\n"); usage.append(" -boolStyle Style for booleans [TRUE_FALSE]\n"); + usage.append(" -blobFormat [base64|hex] Format of blobs: base-64 encoded or hex string [base64]\n"); usage.append(" -numThreads Number of concurrent threads to unload [5]\n"); usage.append(" -beginToken Begin token [none]\n"); usage.append(" -endToken End token [none]\n"); @@ -284,6 +283,13 @@ private boolean parseArgs(String[] args) return false; } } + if (null != (tkey = amap.remove("-blobFormat"))) { + blobFormat = ByteBufferParser.getBlobFormat(tkey); + if (null == blobFormat) { + System.err.println("Bad blob format. Options are: " + ByteBufferParser.getOptions()); + return false; + } + } if (null != (tkey = amap.remove("-numThreads"))) numThreads = Integer.parseInt(tkey); if (null != (tkey = amap.remove("-beginToken"))) beginToken = tkey; if (null != (tkey = amap.remove("-endToken"))) endToken = tkey; @@ -392,7 +398,9 @@ public boolean run(String[] args) nullString, dateFormatString, localDateFormatString, - boolStyle, locale, + boolStyle, + blobFormat, + locale, pstream, beginToken, endToken, session, @@ -443,7 +451,9 @@ public boolean run(String[] args) nullString, dateFormatString, localDateFormatString, - boolStyle, locale, + boolStyle, + blobFormat, + locale, pstream, tBeginString, tEndString, session, @@ -486,6 +496,7 @@ class ThreadExecute implements Callable { private String format = "delim"; private Locale locale = null; private BooleanParser.BoolStyle boolStyle = null; + private ByteBufferParser.BlobFormat blobFormat = null; private String nullString = null; private String delimiter = null; @@ -503,7 +514,8 @@ public ThreadExecute(String inCqlSchema, String inDelimiter, String inNullString, String inDateFormatString, String inLocalDateFormatString, - BooleanParser.BoolStyle inBoolStyle, + BooleanParser.BoolStyle inBoolStyle, + ByteBufferParser.BlobFormat inBlobFormat, Locale inLocale, PrintStream inWriter, String inBeginToken, String inEndToken, @@ -516,6 +528,7 @@ public ThreadExecute(String inCqlSchema, String inDelimiter, dateFormatString = inDateFormatString; localDateFormatString = inLocalDateFormatString; boolStyle = inBoolStyle; + blobFormat = inBlobFormat; locale = inLocale; beginToken = inBeginToken; endToken = inEndToken; @@ -560,7 +573,8 @@ private String getPartitionKey(CqlDelimParser cdp, Session session) { private boolean setup() throws IOException, ParseException { cdp = new CqlDelimParser(cqlSchema, delimiter, 4096, nullString, null, dateFormatString, localDateFormatString, - boolStyle, locale, null, session, false, -1); + boolStyle, blobFormat, locale, null, + session, false, -1); String select = cdp.generateSelect(); String partitionKey = getPartitionKey(cdp, session); if (null != beginToken) { diff --git a/src/main/java/com/datastax/loader/parser/ByteBufferParser.java b/src/main/java/com/datastax/loader/parser/ByteBufferParser.java index 4373612..371124b 100644 --- a/src/main/java/com/datastax/loader/parser/ByteBufferParser.java +++ b/src/main/java/com/datastax/loader/parser/ByteBufferParser.java @@ -17,20 +17,71 @@ import java.nio.ByteBuffer; import javax.xml.bind.DatatypeConverter; -import com.datastax.driver.core.Row; -import com.datastax.driver.core.exceptions.InvalidTypeException; +import com.datastax.driver.core.utils.Bytes; // ByteBuffer parser public class ByteBufferParser extends AbstractParser { + public static enum BlobFormat { + Base64("base64"), + HexString("hex"); + + private String formatStr; + + BlobFormat(String inFormatStr) { + formatStr = inFormatStr; + } + + public String getFormat() { + return formatStr; + } + } + + public static BlobFormat getBlobFormat(String instr) { + for (BlobFormat bs : BlobFormat.values()) { + if (bs.getFormat().equalsIgnoreCase(instr)) { + return bs; + } + } + return null; + } + + public static String getOptions() { + String ret = "'" + BlobFormat.Base64.getFormat() + "'"; + ret = ret + ", '" + BlobFormat.HexString.getFormat() + "'"; + return ret; + } + + private BlobFormat format; + + public ByteBufferParser(BlobFormat inFormat) { + if (null == inFormat) + inFormat = BlobFormat.Base64; + this.format = inFormat; + } + public ByteBuffer parseIt(String toparse) { if (null == toparse) return null; - byte[] barry = DatatypeConverter.parseBase64Binary(toparse); - return ByteBuffer.wrap(barry); + + switch (format) { + case HexString: + ByteBuffer bb = Bytes.fromHexString(toparse); + return bb; + + default: + byte[] barry = DatatypeConverter.parseBase64Binary(toparse); + return ByteBuffer.wrap(barry); + } } public String format(Object o) { ByteBuffer v = (ByteBuffer)o; - return DatatypeConverter.printBase64Binary(v.array()); + switch (format) { + case HexString: + return Bytes.toHexString(v); + + default: + return DatatypeConverter.printBase64Binary(v.array()); + } } }