diff --git a/pom.xml b/pom.xml
index a7806470a1..b4acbb30f7 100644
--- a/pom.xml
+++ b/pom.xml
@@ -69,15 +69,12 @@
2.7.10
1.7.5
4.4.1
- 1.2.1
2.7.1
- 0.94.25
- 0.96.2
- ${hbase096.core.version}-hadoop2
0.98.2
${hbase098.core.version}-hadoop2
1.0.2
${hbase100.core.version}
+ 1.1.1
1.9.2
2.3.0
-
+ titan-dist-hadoop-2
+
${project.groupId}
diff --git a/titan-dist/titan-dist-hadoop-1/pom.xml b/titan-dist/titan-dist-hadoop-1/pom.xml
deleted file mode 100644
index b91030b798..0000000000
--- a/titan-dist/titan-dist-hadoop-1/pom.xml
+++ /dev/null
@@ -1,214 +0,0 @@
-
- 4.0.0
-
- com.thinkaurelius.titan
- titan-dist
- 1.0.1-SNAPSHOT
- ../pom.xml
-
- pom
- titan-dist-hadoop-1
- Titan-Dist: Archive with Hadoop 1
- http://thinkaurelius.github.com/titan/
-
-
- hadoop1
- ${project.parent.basedir}/src/assembly/descriptor
- ${project.parent.basedir}/src/assembly/static
- ${project.parent.basedir}/src/assembly/resources
- ${project.parent.build.directory}/cfilter
- ${project.parent.parent.basedir}
-
-
-
-
- ${project.groupId}
- titan-all
- ${project.version}
- pom
-
-
- org.apache.hadoop
- hadoop-client
-
-
- org.apache.hbase
- hbase-client
-
-
-
- ${project.groupId}
- titan-solr
-
-
-
-
- org.apache.hbase
- hbase-client
- ${hbase098.core.version}-hadoop1
-
-
-
- org.jruby
- jruby-complete
-
-
-
-
- org.apache.hadoop
- hadoop-client
- ${hadoop1.version}
-
-
-
-
-
-
- maven-assembly-plugin
-
-
-
- org.codehaus.mojo
- exec-maven-plugin
-
-
- generate-titan-gremlin-imports
- generate-resources
-
-
-
-
-
- org.codehaus.mojo
- properties-maven-plugin
-
-
-
- maven-gpg-plugin
-
-
-
-
-
-
- aurelius-release
-
-
-
- maven-failsafe-plugin
-
-
-
- org.codehaus.mojo
- wagon-maven-plugin
-
-
-
- maven-resources-plugin
-
-
- filter-cassandra-murmur-config
- process-test-resources
-
-
- filter-cassandra-bop-config
- process-test-resources
-
-
- filter-expect-scripts
- process-test-resources
-
-
- filter-test-resources
- process-test-resources
-
-
- copy-test-cfiles
- process-test-resources
-
-
-
-
-
-
-
-
-
-
-
- dev-install-hadoop1
-
-
-
- dev.hadoop
- 1
-
-
-
-
-
-
- maven-clean-plugin
-
-
- clean-dev-dirs-hadoop
- clean
-
-
-
-
-
- maven-assembly-plugin
-
-
- install-dev-dirs-hadoop
- install
-
-
-
-
-
-
-
-
- dev-install-hadoop1-by-default
-
-
-
-
- !dev.hadoop
-
-
-
-
-
-
- maven-clean-plugin
-
-
- clean-dev-dirs-hadoop
- clean
-
-
-
-
-
- maven-assembly-plugin
-
-
- install-dev-dirs-hadoop
- install
-
-
-
-
-
-
-
-
diff --git a/titan-dist/titan-dist-hadoop-1/src/assembly/descriptor/archive.xml b/titan-dist/titan-dist-hadoop-1/src/assembly/descriptor/archive.xml
deleted file mode 100644
index 399f49ea64..0000000000
--- a/titan-dist/titan-dist-hadoop-1/src/assembly/descriptor/archive.xml
+++ /dev/null
@@ -1,18 +0,0 @@
-
-
- titan-${project.version}-${hadoop.version.tag}
- titan-${project.version}-${hadoop.version.tag}
-
-
- ${assembly.descriptor.dir}/archive.xml
-
-
-
-
- src/assembly/static
- /
- false
-
-
-
diff --git a/titan-dist/titan-dist-hadoop-2/pom.xml b/titan-dist/titan-dist-hadoop-2/pom.xml
index fa5fdaa041..88841f8398 100644
--- a/titan-dist/titan-dist-hadoop-2/pom.xml
+++ b/titan-dist/titan-dist-hadoop-2/pom.xml
@@ -3,7 +3,7 @@
com.thinkaurelius.titan
titan-dist
- 0.9.0-SNAPSHOT
+ 1.0.1-SNAPSHOT
../pom.xml
pom
@@ -24,13 +24,19 @@
org.apache.hbase
hbase-client
- ${hbase098.version}
+ ${hbase111.version}
org.apache.hbase
hbase-server
+ ${hbase111.version}
- ${hbase098.version}
+
+
+ com.lmax
+ disruptor
+
+
@@ -151,7 +157,7 @@
-classpath
com.thinkaurelius.titan.example.GraphOfTheGodsFactory
- ${project.build.directory}/conf/titan-berkeleydb-es.properties
+ ${project.build.directory}/conf/titan-berkeleyje-es.properties
diff --git a/titan-hadoop-parent/pom.xml b/titan-hadoop-parent/pom.xml
index 363c4f0c50..3456d7bb2e 100644
--- a/titan-hadoop-parent/pom.xml
+++ b/titan-hadoop-parent/pom.xml
@@ -20,7 +20,6 @@
titan-hadoop-core
- titan-hadoop-1
titan-hadoop-2
titan-hadoop
diff --git a/titan-hadoop-parent/titan-hadoop-1/pom.xml b/titan-hadoop-parent/titan-hadoop-1/pom.xml
deleted file mode 100644
index 69cb61d6b9..0000000000
--- a/titan-hadoop-parent/titan-hadoop-1/pom.xml
+++ /dev/null
@@ -1,112 +0,0 @@
-
- 4.0.0
-
- com.thinkaurelius.titan
- titan-hadoop-parent
- 1.0.1-SNAPSHOT
- ../pom.xml
-
- titan-hadoop-1
- Titan-Hadoop: 1.x Compatibility Shim
- http://thinkaurelius.github.com/titan/
-
-
- ${basedir}/../..
-
-
-
-
- org.apache.hadoop
- hadoop-core
- ${hadoop1.version}
-
-
- javax.servlet
- jsp-api
-
-
- javax.servlet
- servlet-api
-
-
- org.mortbay.jetty
- servlet-api-2.5
-
-
- org.mortbay.jetty
- servlet-api
-
-
-
-
- ${project.groupId}
- titan-hadoop-core
- ${project.version}
-
-
- ${project.groupId}
- titan-hadoop-core
- ${project.version}
- tests
- test
-
-
- ${project.groupId}
- titan-hadoop-core
- ${project.version}
- shared-resources
- test
- true
-
-
- org.apache.hbase
- hbase-client
- ${hbase098.core.version}-hadoop1
- true
- test
-
-
-
- org.jruby
- jruby-complete
-
-
-
-
- org.apache.mrunit
- mrunit
- ${mrunit.version}
- hadoop1
-
-
-
- ${project.groupId}
- titan-es
- ${project.version}
- tests
- test
-
-
-
-
-
-
- maven-assembly-plugin
-
-
- maven-dependency-plugin
-
-
- maven-failsafe-plugin
-
-
- maven-surefire-plugin
-
-
-
-
diff --git a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/DistCacheConfigurer.java b/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/DistCacheConfigurer.java
deleted file mode 100644
index 4432044208..0000000000
--- a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/DistCacheConfigurer.java
+++ /dev/null
@@ -1,36 +0,0 @@
-package com.thinkaurelius.titan.hadoop.compat.h1;
-
-import com.thinkaurelius.titan.hadoop.config.job.AbstractDistCacheConfigurer;
-import com.thinkaurelius.titan.hadoop.config.job.JobClasspathConfigurer;
-import org.apache.hadoop.conf.Configuration;
-import org.apache.hadoop.filecache.DistributedCache;
-import org.apache.hadoop.fs.FileSystem;
-import org.apache.hadoop.fs.Path;
-import org.apache.hadoop.mapreduce.Job;
-
-import java.io.IOException;
-
-public class DistCacheConfigurer extends AbstractDistCacheConfigurer implements JobClasspathConfigurer {
-
- public DistCacheConfigurer(String mapredJarFilename) {
- super(mapredJarFilename);
- }
-
- @Override
- public void configure(Job job) throws IOException {
-
- for (Path p : getLocalPaths()) {
- Configuration conf = job.getConfiguration();
- FileSystem jobFS = FileSystem.get(conf);
- FileSystem localFS = FileSystem.getLocal(conf);
- Path stagedPath = uploadFileIfNecessary(localFS, p, jobFS);
- DistributedCache.addFileToClassPath(stagedPath, conf, jobFS);
- }
-
- // We don't really need to set a mapred job jar here,
- // but doing so suppresses a warning
- String mj = getMapredJar();
- if (null != mj)
- job.getConfiguration().set(Hadoop1Compat.CFG_JOB_JAR, mj);
- }
-}
diff --git a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/Hadoop1Compat.java b/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/Hadoop1Compat.java
deleted file mode 100644
index b47025a882..0000000000
--- a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/Hadoop1Compat.java
+++ /dev/null
@@ -1,89 +0,0 @@
-package com.thinkaurelius.titan.hadoop.compat.h1;
-
-import com.thinkaurelius.titan.core.TitanException;
-import com.thinkaurelius.titan.diskstorage.keycolumnvalue.scan.ScanMetrics;
-import com.thinkaurelius.titan.graphdb.configuration.TitanConstants;
-import com.thinkaurelius.titan.hadoop.config.job.JobClasspathConfigurer;
-import com.thinkaurelius.titan.hadoop.formats.cassandra.CassandraBinaryInputFormat;
-import com.thinkaurelius.titan.hadoop.scan.HadoopVertexScanMapper;
-import org.apache.hadoop.conf.Configuration;
-import org.apache.hadoop.io.NullWritable;
-import org.apache.hadoop.mapreduce.*;
-import org.apache.hadoop.mapreduce.lib.output.NullOutputFormat;
-
-import com.thinkaurelius.titan.hadoop.compat.HadoopCompat;
-
-import java.io.IOException;
-
-public class Hadoop1Compat implements HadoopCompat {
-
- static final String CFG_SPECULATIVE_MAPS = "mapred.map.tasks.speculative.execution";
- static final String CFG_SPECULATIVE_REDUCES = "mapred.reduce.tasks.speculative.execution";
- static final String CFG_JOB_JAR = "mapred.jar";
-
- @Override
- public TaskAttemptContext newTask(Configuration c, TaskAttemptID t) {
- return new TaskAttemptContext(c, t);
- }
-
- @Override
- public String getSpeculativeMapConfigKey() {
- return CFG_SPECULATIVE_MAPS;
- }
-
- @Override
- public String getSpeculativeReduceConfigKey() {
- return CFG_SPECULATIVE_REDUCES;
- }
-
- @Override
- public String getMapredJarConfigKey() {
- return CFG_JOB_JAR;
- }
-
- @Override
- public long getContextCounter(TaskInputOutputContext context, String group, String name) {
- return context.getCounter(group, name).getValue();
- }
-
- @Override
- public void incrementContextCounter(TaskInputOutputContext context,
- String group, String name, long incr) {
- context.getCounter(group, name).increment(incr);
- }
-
- @Override
- public Configuration getContextConfiguration(TaskAttemptContext context) {
- return context.getConfiguration();
- }
-
- @Override
- public JobClasspathConfigurer newMapredJarConfigurer(String mapredJarPath) {
- return new MapredJarConfigurer(mapredJarPath);
- }
-
- @Override
- public JobClasspathConfigurer newDistCacheConfigurer() {
- return new DistCacheConfigurer("titan-hadoop-core-" + TitanConstants.VERSION + ".jar");
- }
-
- @Override
- public Configuration getJobContextConfiguration(JobContext context) {
- return context.getConfiguration();
- }
-
- @Override
- public Configuration newImmutableConfiguration(Configuration base) {
- return new ImmutableConfiguration(base);
- }
-
- @Override
- public ScanMetrics getMetrics(Counters c) {
- return new Hadoop1CountersScanMetrics(c);
- }
-
- @Override
- public String getJobFailureString(Job j) {
- return j.toString();
- }
-}
diff --git a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/Hadoop1CountersScanMetrics.java b/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/Hadoop1CountersScanMetrics.java
deleted file mode 100644
index a9c236d3fc..0000000000
--- a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/Hadoop1CountersScanMetrics.java
+++ /dev/null
@@ -1,39 +0,0 @@
-package com.thinkaurelius.titan.hadoop.compat.h1;
-
-import com.thinkaurelius.titan.diskstorage.keycolumnvalue.scan.ScanMetrics;
-import com.thinkaurelius.titan.hadoop.scan.HadoopContextScanMetrics;
-import org.apache.hadoop.mapreduce.Counters;
-
-public class Hadoop1CountersScanMetrics implements ScanMetrics {
-
- private final Counters counters;
-
- public Hadoop1CountersScanMetrics(Counters counters) {
- this.counters = counters;
- }
-
- @Override
- public long getCustom(String metric) {
- return counters.getGroup(HadoopContextScanMetrics.CUSTOM_COUNTER_GROUP).findCounter(metric).getValue();
- }
-
- @Override
- public void incrementCustom(String metric, long delta) {
- counters.getGroup(HadoopContextScanMetrics.CUSTOM_COUNTER_GROUP).findCounter(metric).increment(delta);
- }
-
- @Override
- public void incrementCustom(String metric) {
- incrementCustom(metric, 1L);
- }
-
- @Override
- public long get(Metric metric) {
- return counters.getGroup(HadoopContextScanMetrics.STANDARD_COUNTER_GROUP).findCounter(metric.name()).getValue();
- }
-
- @Override
- public void increment(Metric metric) {
- counters.getGroup(HadoopContextScanMetrics.STANDARD_COUNTER_GROUP).findCounter(metric.name()).increment(1L);
- }
-}
diff --git a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/ImmutableConfiguration.java b/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/ImmutableConfiguration.java
deleted file mode 100644
index b85647f018..0000000000
--- a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/ImmutableConfiguration.java
+++ /dev/null
@@ -1,283 +0,0 @@
-package com.thinkaurelius.titan.hadoop.compat.h1;
-
-import org.apache.hadoop.conf.Configuration;
-import org.apache.hadoop.fs.Path;
-
-import java.io.*;
-import java.net.URL;
-import java.util.Collection;
-import java.util.Iterator;
-import java.util.List;
-import java.util.Map;
-
-public class ImmutableConfiguration extends Configuration {
-
- private final Configuration encapsulated;
-
- public ImmutableConfiguration(final Configuration encapsulated) {
- this.encapsulated = encapsulated;
- }
-
- public static void addDefaultResource(String name) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public void addResource(String name) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public void addResource(URL url) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public void addResource(Path file) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public void addResource(InputStream in) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public void reloadConfiguration() {
- //throw new UnsupportedOperationException("This configuration instance is immutable");
- encapsulated.reloadConfiguration(); // allowed to simplify testing
- }
-
- @Override
- public String get(String name) {
- return encapsulated.get(name);
- }
-
- @Override
- public String getRaw(String name) {
- return encapsulated.getRaw(name);
- }
-
- @Override
- public void set(String name, String value) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public void unset(String name) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public void setIfUnset(String name, String value) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public String get(String name, String defaultValue) {
- return encapsulated.get(name, defaultValue);
- }
-
- @Override
- public int getInt(String name, int defaultValue) {
- return encapsulated.getInt(name, defaultValue);
- }
-
- @Override
- public void setInt(String name, int value) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public long getLong(String name, long defaultValue) {
- return encapsulated.getLong(name, defaultValue);
- }
-
- @Override
- public void setLong(String name, long value) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public float getFloat(String name, float defaultValue) {
- return encapsulated.getFloat(name, defaultValue);
- }
-
- @Override
- public void setFloat(String name, float value) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public boolean getBoolean(String name, boolean defaultValue) {
- return encapsulated.getBoolean(name, defaultValue);
- }
-
- @Override
- public void setBoolean(String name, boolean value) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public void setBooleanIfUnset(String name, boolean value) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public > void setEnum(String name, T value) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public > T getEnum(String name, T defaultValue) {
- return encapsulated.getEnum(name, defaultValue);
- }
-
- @Override
- public IntegerRanges getRange(String name, String defaultValue) {
- return encapsulated.getRange(name, defaultValue);
- }
-
- @Override
- public Collection getStringCollection(String name) {
- return encapsulated.getStringCollection(name);
- }
-
- @Override
- public String[] getStrings(String name) {
- return encapsulated.getStrings(name);
- }
-
- @Override
- public String[] getStrings(String name, String... defaultValue) {
- return encapsulated.getStrings(name, defaultValue);
- }
-
- @Override
- public void setStrings(String name, String... values) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public Class> getClassByName(String name) throws ClassNotFoundException {
- return encapsulated.getClassByName(name);
- }
-
- @Override
- public Class>[] getClasses(String name, Class>... defaultValue) {
- return encapsulated.getClasses(name, defaultValue);
- }
-
- @Override
- public Class> getClass(String name, Class> defaultValue) {
- return encapsulated.getClass(name, defaultValue);
- }
-
- @Override
- public Class extends U> getClass(String name, Class extends U> defaultValue, Class xface) {
- return encapsulated.getClass(name, defaultValue, xface);
- }
-
- @Override
- public List getInstances(String name, Class xface) {
- return encapsulated.getInstances(name, xface);
- }
-
- @Override
- public void setClass(String name, Class> theClass, Class> xface) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public Path getLocalPath(String dirsProp, String path) throws IOException {
- return encapsulated.getLocalPath(dirsProp, path);
- }
-
- @Override
- public File getFile(String dirsProp, String path) throws IOException {
- return encapsulated.getFile(dirsProp, path);
- }
-
- @Override
- public URL getResource(String name) {
- return encapsulated.getResource(name);
- }
-
- @Override
- public InputStream getConfResourceAsInputStream(String name) {
- return encapsulated.getConfResourceAsInputStream(name);
- }
-
- @Override
- public Reader getConfResourceAsReader(String name) {
- return encapsulated.getConfResourceAsReader(name);
- }
-
- @Override
- public int size() {
- return encapsulated.size();
- }
-
- @Override
- public void clear() {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public Iterator> iterator() {
- return encapsulated.iterator();
- }
-
- @Override
- public void writeXml(OutputStream out) throws IOException {
- encapsulated.writeXml(out);
- }
-
- @Override
- public void writeXml(Writer out) throws IOException {
- encapsulated.writeXml(out);
- }
-
- public static void dumpConfiguration(Configuration config, Writer out) throws IOException {
- Configuration.dumpConfiguration(config, out);
- }
-
- @Override
- public ClassLoader getClassLoader() {
- return encapsulated.getClassLoader();
- }
-
- @Override
- public void setClassLoader(ClassLoader classLoader) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- @Override
- public String toString() {
- return encapsulated.toString();
- }
-
- @Override
- public void setQuietMode(boolean quietmode) {
- throw new UnsupportedOperationException("This configuration instance is immutable");
- }
-
- public static void main(String[] args) throws Exception {
- Configuration.main(args);
- }
-
- @Override
- public void readFields(DataInput in) throws IOException {
- encapsulated.readFields(in);
- }
-
- @Override
- public void write(DataOutput out) throws IOException {
- encapsulated.write(out);
- }
-
- @Override
- public Map getValByRegex(String regex) {
- return encapsulated.getValByRegex(regex);
- }
-}
diff --git a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/MapredJarConfigurer.java b/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/MapredJarConfigurer.java
deleted file mode 100644
index 2be1c5f013..0000000000
--- a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/compat/h1/MapredJarConfigurer.java
+++ /dev/null
@@ -1,20 +0,0 @@
-package com.thinkaurelius.titan.hadoop.compat.h1;
-
-import com.thinkaurelius.titan.hadoop.config.job.JobClasspathConfigurer;
-import org.apache.hadoop.mapreduce.Job;
-
-import java.io.IOException;
-
-public class MapredJarConfigurer implements JobClasspathConfigurer {
-
- private final String mapredJar;
-
- public MapredJarConfigurer(String mapredJar) {
- this.mapredJar = mapredJar;
- }
-
- @Override
- public void configure(Job job) throws IOException {
- job.getConfiguration().set(Hadoop1Compat.CFG_JOB_JAR, mapredJar);
- }
-}
diff --git a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/formats/TitanH1OutputCommitter.java b/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/formats/TitanH1OutputCommitter.java
deleted file mode 100644
index 2134ed222d..0000000000
--- a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/formats/TitanH1OutputCommitter.java
+++ /dev/null
@@ -1,40 +0,0 @@
-package com.thinkaurelius.titan.hadoop.formats;
-
-import org.apache.hadoop.mapreduce.JobContext;
-import org.apache.hadoop.mapreduce.OutputCommitter;
-import org.apache.hadoop.mapreduce.TaskAttemptContext;
-
-import java.io.IOException;
-
-public class TitanH1OutputCommitter extends OutputCommitter {
- private final TitanH1OutputFormat tof;
-
- public TitanH1OutputCommitter(TitanH1OutputFormat tof) {
- this.tof = tof;
- }
-
- @Override
- public void setupJob(JobContext jobContext) throws IOException {
-
- }
-
- @Override
- public void setupTask(TaskAttemptContext taskContext) throws IOException {
-
- }
-
- @Override
- public boolean needsTaskCommit(TaskAttemptContext taskContext) throws IOException {
- return tof.hasModifications(taskContext.getTaskAttemptID());
- }
-
- @Override
- public void commitTask(TaskAttemptContext taskContext) throws IOException {
- tof.commit(taskContext.getTaskAttemptID());
- }
-
- @Override
- public void abortTask(TaskAttemptContext taskContext) throws IOException {
- tof.abort(taskContext.getTaskAttemptID());
- }
-}
diff --git a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/formats/TitanH1OutputFormat.java b/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/formats/TitanH1OutputFormat.java
deleted file mode 100644
index a898ed5600..0000000000
--- a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/formats/TitanH1OutputFormat.java
+++ /dev/null
@@ -1,103 +0,0 @@
-package com.thinkaurelius.titan.hadoop.formats;
-
-import com.google.common.base.Joiner;
-import com.google.common.collect.ImmutableSet;
-import com.thinkaurelius.titan.core.TitanFactory;
-import com.thinkaurelius.titan.graphdb.database.StandardTitanGraph;
-import com.thinkaurelius.titan.graphdb.transaction.StandardTitanTx;
-import com.thinkaurelius.titan.hadoop.config.ModifiableHadoopConfiguration;
-import com.thinkaurelius.titan.hadoop.config.TitanHadoopConfiguration;
-import org.apache.hadoop.conf.Configuration;
-import org.apache.hadoop.io.NullWritable;
-import org.apache.hadoop.mapreduce.JobContext;
-import org.apache.hadoop.mapreduce.OutputCommitter;
-import org.apache.hadoop.mapreduce.OutputFormat;
-import org.apache.hadoop.mapreduce.RecordWriter;
-import org.apache.hadoop.mapreduce.TaskAttemptContext;
-import org.apache.hadoop.mapreduce.TaskAttemptID;
-import org.apache.tinkerpop.gremlin.hadoop.structure.io.VertexWritable;
-import org.apache.tinkerpop.gremlin.hadoop.structure.util.ConfUtil;
-import org.apache.tinkerpop.gremlin.process.computer.VertexProgram;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.io.IOException;
-import java.util.Set;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.ConcurrentMap;
-
-public class TitanH1OutputFormat extends OutputFormat {
-
- private static final Logger log = LoggerFactory.getLogger(TitanH1OutputFormat.class);
-
- private final ConcurrentMap transactions = new ConcurrentHashMap<>();
-
- private StandardTitanGraph graph;
-
- private Set persistableKeys;
-
- @Override
- public RecordWriter getRecordWriter(TaskAttemptContext taskAttemptContext) throws IOException, InterruptedException {
-
- synchronized (this) {
- if (null == graph) {
- Configuration hadoopConf = taskAttemptContext.getConfiguration();
- ModifiableHadoopConfiguration mhc =
- ModifiableHadoopConfiguration.of(TitanHadoopConfiguration.MAPRED_NS, hadoopConf);
- graph = (StandardTitanGraph) TitanFactory.open(mhc.getTitanGraphConf());
- }
- }
-
- // Special case for a TP3 vertex program: persist only those properties whose keys are
- // returned by VertexProgram.getComputeKeys()
- if (null == persistableKeys) {
- try {
- persistableKeys = VertexProgram.createVertexProgram(graph,
- ConfUtil.makeApacheConfiguration(taskAttemptContext.getConfiguration())).getElementComputeKeys();
- log.debug("Set persistableKeys={}", Joiner.on(",").join(persistableKeys));
- } catch (Exception e) {
- log.debug("Unable to detect or instantiate vertex program", e);
- persistableKeys = ImmutableSet.of();
- }
- }
-
- StandardTitanTx tx = transactions.computeIfAbsent(taskAttemptContext.getTaskAttemptID(),
- id -> (StandardTitanTx)graph.newTransaction());
- return new TitanH1RecordWriter(taskAttemptContext, tx, persistableKeys);
- }
-
- @Override
- public void checkOutputSpecs(JobContext jobContext) throws IOException, InterruptedException {
- // TODO check output configuration for minimum set of keys here?
- }
-
- @Override
- public OutputCommitter getOutputCommitter(TaskAttemptContext taskAttemptContext) throws IOException,
- InterruptedException {
- return new TitanH1OutputCommitter(this);
- }
-
- void commit(TaskAttemptID id) {
- StandardTitanTx tx = transactions.remove(id);
- if (null == tx) {
- log.warn("Detected concurrency in task commit");
- return;
- }
- tx.commit();
- }
-
- void abort(TaskAttemptID id) {
- StandardTitanTx tx = transactions.remove(id);
- if (null == tx) {
- log.warn("Detected concurrency in task abort");
- return;
- }
- tx.rollback();
- }
-
- boolean hasModifications(TaskAttemptID id) {
- StandardTitanTx tx = transactions.get(id);
- // if tx is null, something is horribly wrong
- return tx.hasModifications();
- }
-}
diff --git a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/formats/TitanH1RecordWriter.java b/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/formats/TitanH1RecordWriter.java
deleted file mode 100644
index b2f4f28e8b..0000000000
--- a/titan-hadoop-parent/titan-hadoop-1/src/main/java/com/thinkaurelius/titan/hadoop/formats/TitanH1RecordWriter.java
+++ /dev/null
@@ -1,53 +0,0 @@
-package com.thinkaurelius.titan.hadoop.formats;
-
-import com.thinkaurelius.titan.graphdb.transaction.StandardTitanTx;
-import org.apache.hadoop.io.NullWritable;
-import org.apache.hadoop.mapreduce.RecordWriter;
-import org.apache.hadoop.mapreduce.TaskAttemptContext;
-import org.apache.tinkerpop.gremlin.hadoop.structure.io.VertexWritable;
-import org.apache.tinkerpop.gremlin.structure.Vertex;
-import org.apache.tinkerpop.gremlin.structure.VertexProperty;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.io.IOException;
-import java.util.Iterator;
-import java.util.Set;
-
-public class TitanH1RecordWriter extends RecordWriter {
-
- private static final Logger log = LoggerFactory.getLogger(TitanH1RecordWriter.class);
-
- private final TaskAttemptContext taskAttemptContext;
- private final StandardTitanTx tx;
- private final Set persistableKeys;
-
- public TitanH1RecordWriter(TaskAttemptContext taskAttemptContext, StandardTitanTx tx, Set persistableKeys) {
- this.taskAttemptContext = taskAttemptContext;
- this.tx = tx;
- this.persistableKeys = persistableKeys;
- }
-
- @Override
- public void write(NullWritable key, VertexWritable value) throws IOException, InterruptedException {
- // TODO tolerate possibility that concurrent OLTP activity has deleted the vertex? maybe configurable...
- Object vertexID = value.get().id();
- Vertex vertex = tx.vertices(vertexID).next();
- Iterator> vpIter = value.get().properties();
- while (vpIter.hasNext()) {
- VertexProperty