Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
@@ -1,8 +1,16 @@
!.gitignore

out/
build/
gradle/
.gradle/
.idea/
*~
*.jar
cassandra-loader
*.BADINSERT
*.BADPARSE
*.LOG
*.iml
/gradlew
/gradlew.bat
2 changes: 1 addition & 1 deletion build.gradle
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
apply plugin: 'java'
apply plugin: 'application'

def versionNum = '0.0.27'
def versionNum = '0.0.28'

allprojects {
tasks.withType(JavaCompile) {
Expand Down
38 changes: 18 additions & 20 deletions src/main/java/com/datastax/loader/CqlDelimLoad.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -169,6 +155,7 @@ private String usage() {
usage.append(" -batchSize <batchSize> Number of INSERTs to batch together [1]\n");
usage.append(" -decimalDelim <decimalDelim> Decimal delimiter [.] Other option is ','\n");
usage.append(" -boolStyle <boolStyleString> 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 <numThreads> Number of concurrent threads (files) to load [num cores]\n");
usage.append(" -queryTimeout <# seconds> Query timeout (in seconds) [2]\n");
usage.append(" -numRetries <numRetries> Number of times to retry the INSERT [1]\n");
Expand All @@ -178,7 +165,7 @@ private String usage() {
usage.append(" -rateFile <filename> Where to print the rate statistics\n");
usage.append(" -successDir <dir> Directory where to move successfully loaded files\n");
usage.append(" -failureDir <dir> 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 <tableName> Table name (when using JSON)\n");
usage.append(" -keyspace <keyspaceName> Keyspace name (when using JSON)\n");
Expand Down Expand Up @@ -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);

Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand Down
16 changes: 11 additions & 5 deletions src/main/java/com/datastax/loader/CqlDelimLoadTask.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -84,6 +85,7 @@ class CqlDelimLoadTask implements Callable<Long> {
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;
Expand All @@ -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,
Expand All @@ -125,6 +128,7 @@ public CqlDelimLoadTask(String inCqlSchema, String inDelimiter,
dateFormatString = inDateFormatString;
localDateFormatString = inLocalDateFormatString;
boolStyle = inBoolStyle;
blobFormat = inBlobFormat;
locale = inLocale;
maxErrors = inMaxErrors;
skipRows = inSkipRows;
Expand Down Expand Up @@ -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();
Expand Down
27 changes: 15 additions & 12 deletions src/main/java/com/datastax/loader/CqlDelimParser.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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);
}
Expand All @@ -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);
}
Expand All @@ -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<DataType.Name, Parser>();
Parser byteParser = new ByteParser(inLocale, bLoader);
Expand All @@ -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);
Expand Down
32 changes: 23 additions & 9 deletions src/main/java/com/datastax/loader/CqlDelimUnload.java
Original file line number Diff line number Diff line change
Expand Up @@ -23,23 +23,19 @@
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;
import java.io.PrintStream;
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;
Expand Down Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -124,6 +122,7 @@ private String usage() {
usage.append(" -consistencyLevel <CL> Consistency level [LOCAL_ONE]\n");
usage.append(" -decimalDelim <decimalDelim> Decimal delimiter [.] Other option is ','\n");
usage.append(" -boolStyle <boolStyleString> 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 <numThreads> Number of concurrent threads to unload [5]\n");
usage.append(" -beginToken <tokenString> Begin token [none]\n");
usage.append(" -endToken <tokenString> End token [none]\n");
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -392,7 +398,9 @@ public boolean run(String[] args)
nullString,
dateFormatString,
localDateFormatString,
boolStyle, locale,
boolStyle,
blobFormat,
locale,
pstream,
beginToken,
endToken, session,
Expand Down Expand Up @@ -443,7 +451,9 @@ public boolean run(String[] args)
nullString,
dateFormatString,
localDateFormatString,
boolStyle, locale,
boolStyle,
blobFormat,
locale,
pstream,
tBeginString,
tEndString, session,
Expand Down Expand Up @@ -486,6 +496,7 @@ class ThreadExecute implements Callable<Long> {
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;

Expand All @@ -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,
Expand All @@ -516,6 +528,7 @@ public ThreadExecute(String inCqlSchema, String inDelimiter,
dateFormatString = inDateFormatString;
localDateFormatString = inLocalDateFormatString;
boolStyle = inBoolStyle;
blobFormat = inBlobFormat;
locale = inLocale;
beginToken = inBeginToken;
endToken = inEndToken;
Expand Down Expand Up @@ -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) {
Expand Down
Loading