diff --git a/examples/hadoop_cql3_word_count/src/WordCount.java b/examples/hadoop_cql3_word_count/src/WordCount.java index 09dd9e47ec..611f9c2911 100644 --- a/examples/hadoop_cql3_word_count/src/WordCount.java +++ b/examples/hadoop_cql3_word_count/src/WordCount.java @@ -21,13 +21,12 @@ import java.nio.ByteBuffer; import java.util.*; import java.util.Map.Entry; -import org.apache.cassandra.thrift.*; -import org.apache.cassandra.hadoop.cql3.ColumnFamilyOutputFormat; +import org.apache.cassandra.hadoop.cql3.CqlConfigHelper; +import org.apache.cassandra.hadoop.cql3.CqlOutputFormat; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.hadoop.cql3.CQLConfigHelper; -import org.apache.cassandra.hadoop.cql3.ColumnFamilyInputFormat; +import org.apache.cassandra.hadoop.cql3.CqlPagingInputFormat; import org.apache.cassandra.hadoop.ConfigHelper; import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.hadoop.conf.Configuration; @@ -38,7 +37,6 @@ import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.Mapper; import org.apache.hadoop.mapreduce.Reducer; -import org.apache.hadoop.mapreduce.Reducer.Context; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import org.apache.hadoop.util.Tool; import org.apache.hadoop.util.ToolRunner; @@ -208,28 +206,28 @@ public class WordCount extends Configured implements Tool job.setOutputKeyClass(Map.class); job.setOutputValueClass(List.class); - job.setOutputFormatClass(ColumnFamilyOutputFormat.class); + job.setOutputFormatClass(CqlOutputFormat.class); ConfigHelper.setOutputColumnFamily(job.getConfiguration(), KEYSPACE, OUTPUT_COLUMN_FAMILY); job.getConfiguration().set(PRIMARY_KEY, "word,sum"); String query = "INSERT INTO " + KEYSPACE + "." + OUTPUT_COLUMN_FAMILY + " (row_id1, row_id2, word, count_num) " + " values (?, ?, ?, ?)"; - CQLConfigHelper.setOutputCql(job.getConfiguration(), query); + CqlConfigHelper.setOutputCql(job.getConfiguration(), query); ConfigHelper.setOutputInitialAddress(job.getConfiguration(), "localhost"); ConfigHelper.setOutputPartitioner(job.getConfiguration(), "Murmur3Partitioner"); } - job.setInputFormatClass(ColumnFamilyInputFormat.class); + job.setInputFormatClass(CqlPagingInputFormat.class); ConfigHelper.setInputRpcPort(job.getConfiguration(), "9160"); ConfigHelper.setInputInitialAddress(job.getConfiguration(), "localhost"); ConfigHelper.setInputColumnFamily(job.getConfiguration(), KEYSPACE, COLUMN_FAMILY); ConfigHelper.setInputPartitioner(job.getConfiguration(), "Murmur3Partitioner"); - CQLConfigHelper.setInputCQLPageRowSize(job.getConfiguration(), "3"); + CqlConfigHelper.setInputCQLPageRowSize(job.getConfiguration(), "3"); //this is the user defined filter clauses, you can comment it out if you want count all titles - CQLConfigHelper.setInputWhereClauses(job.getConfiguration(), "title='A'"); + CqlConfigHelper.setInputWhereClauses(job.getConfiguration(), "title='A'"); job.waitForCompletion(true); return 0; } diff --git a/examples/hadoop_cql3_word_count/src/WordCountCounters.java b/examples/hadoop_cql3_word_count/src/WordCountCounters.java index 1cf5539460..8454b7002a 100644 --- a/examples/hadoop_cql3_word_count/src/WordCountCounters.java +++ b/examples/hadoop_cql3_word_count/src/WordCountCounters.java @@ -24,23 +24,20 @@ import java.util.*; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.hadoop.cql3.CqlConfigHelper; +import org.apache.cassandra.hadoop.cql3.CqlPagingInputFormat; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.conf.Configured; import org.apache.hadoop.fs.Path; -import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.Mapper; -import org.apache.hadoop.mapreduce.Mapper.Context; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import org.apache.hadoop.util.Tool; import org.apache.hadoop.util.ToolRunner; -import org.apache.cassandra.hadoop.cql3.ColumnFamilyInputFormat; import org.apache.cassandra.hadoop.ConfigHelper; -import org.apache.cassandra.hadoop.cql3.CQLConfigHelper; -import org.apache.cassandra.thrift.*; import org.apache.cassandra.utils.ByteBufferUtil; @@ -107,14 +104,14 @@ public class WordCountCounters extends Configured implements Tool job.setOutputValueClass(LongWritable.class); FileOutputFormat.setOutputPath(job, new Path(OUTPUT_PATH_PREFIX)); - job.setInputFormatClass(ColumnFamilyInputFormat.class); + job.setInputFormatClass(CqlPagingInputFormat.class); ConfigHelper.setInputRpcPort(job.getConfiguration(), "9160"); ConfigHelper.setInputInitialAddress(job.getConfiguration(), "localhost"); ConfigHelper.setInputPartitioner(job.getConfiguration(), "Murmur3Partitioner"); ConfigHelper.setInputColumnFamily(job.getConfiguration(), WordCount.KEYSPACE, WordCount.OUTPUT_COLUMN_FAMILY); - CQLConfigHelper.setInputCQLPageRowSize(job.getConfiguration(), "3"); + CqlConfigHelper.setInputCQLPageRowSize(job.getConfiguration(), "3"); job.waitForCompletion(true); return 0; diff --git a/src/java/org/apache/cassandra/hadoop/cql3/CQLConfigHelper.java b/src/java/org/apache/cassandra/hadoop/cql3/CQLConfigHelper.java index 66bcfdb6eb..cb61d054e0 100644 --- a/src/java/org/apache/cassandra/hadoop/cql3/CQLConfigHelper.java +++ b/src/java/org/apache/cassandra/hadoop/cql3/CQLConfigHelper.java @@ -21,7 +21,7 @@ package org.apache.cassandra.hadoop.cql3; */ import org.apache.hadoop.conf.Configuration; -public class CQLConfigHelper +public class CqlConfigHelper { private static final String INPUT_CQL_COLUMNS_CONFIG = "cassandra.input.columnfamily.columns"; // separate by colon , private static final String INPUT_CQL_PAGE_ROW_SIZE_CONFIG = "cassandra.input.page.row.size"; diff --git a/src/java/org/apache/cassandra/hadoop/cql3/ColumnFamilyOutputFormat.java b/src/java/org/apache/cassandra/hadoop/cql3/CqlOutputFormat.java similarity index 79% rename from src/java/org/apache/cassandra/hadoop/cql3/ColumnFamilyOutputFormat.java rename to src/java/org/apache/cassandra/hadoop/cql3/CqlOutputFormat.java index 3f6e2afafd..d1e93c805a 100644 --- a/src/java/org/apache/cassandra/hadoop/cql3/ColumnFamilyOutputFormat.java +++ b/src/java/org/apache/cassandra/hadoop/cql3/CqlOutputFormat.java @@ -38,7 +38,7 @@ import org.apache.hadoop.mapreduce.*; *
* As is the case with the {@link ColumnFamilyInputFormat}, you need to set the * prepared statement in your - * Hadoop job Configuration. The {@link CQLConfigHelper} class, through its + * Hadoop job Configuration. The {@link CqlConfigHelper} class, through its * {@link ConfigHelper#setOutputPreparedStatement} method, is provided to make this * simple. * you need to set the Keyspace. The {@link ConfigHelper} class, through its @@ -54,13 +54,13 @@ import org.apache.hadoop.mapreduce.*; * to Cassandra. *
*/ -public class ColumnFamilyOutputFormat extends AbstractColumnFamilyOutputFormat