From a0047797638b9c0d30ebf225e8a00a029e05be0a Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 6 Jun 2013 17:00:08 -0500 Subject: [PATCH] rename --- .../hadoop_cql3_word_count/src/WordCount.java | 18 ++++++++---------- .../src/WordCountCounters.java | 11 ++++------- .../cassandra/hadoop/cql3/CQLConfigHelper.java | 2 +- ...yOutputFormat.java => CqlOutputFormat.java} | 12 ++++++------ ...utFormat.java => CqlPagingInputFormat.java} | 6 +++--- ...dReader.java => CqlPagingRecordReader.java} | 14 +++++++------- ...yRecordWriter.java => CqlRecordWriter.java} | 18 +++++++++--------- 7 files changed, 38 insertions(+), 43 deletions(-) rename src/java/org/apache/cassandra/hadoop/cql3/{ColumnFamilyOutputFormat.java => CqlOutputFormat.java} (79%) rename src/java/org/apache/cassandra/hadoop/cql3/{ColumnFamilyInputFormat.java => CqlPagingInputFormat.java} (92%) rename src/java/org/apache/cassandra/hadoop/cql3/{ColumnFamilyRecordReader.java => CqlPagingRecordReader.java} (98%) rename src/java/org/apache/cassandra/hadoop/cql3/{ColumnFamilyRecordWriter.java => CqlRecordWriter.java} (95%) 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, List> +public class CqlOutputFormat extends AbstractColumnFamilyOutputFormat, List> { /** Fills the deprecated OutputFormat interface for streaming. */ @Deprecated - public ColumnFamilyRecordWriter getRecordWriter(org.apache.hadoop.fs.FileSystem filesystem, org.apache.hadoop.mapred.JobConf job, String name, org.apache.hadoop.util.Progressable progress) throws IOException + public CqlRecordWriter getRecordWriter(org.apache.hadoop.fs.FileSystem filesystem, org.apache.hadoop.mapred.JobConf job, String name, org.apache.hadoop.util.Progressable progress) throws IOException { - return new ColumnFamilyRecordWriter(job, new Progressable(progress)); + return new CqlRecordWriter(job, new Progressable(progress)); } /** @@ -71,8 +71,8 @@ public class ColumnFamilyOutputFormat extends AbstractColumnFamilyOutputFormat, Map> +public class CqlPagingInputFormat extends AbstractColumnFamilyInputFormat, Map> { public RecordReader, Map> getRecordReader(InputSplit split, JobConf jobConf, final Reporter reporter) throws IOException @@ -67,7 +67,7 @@ public class ColumnFamilyInputFormat extends AbstractColumnFamilyInputFormat * Map as column name to columns mappings */ -public class ColumnFamilyRecordReader extends RecordReader, Map> +public class CqlPagingRecordReader extends RecordReader, Map> implements org.apache.hadoop.mapred.RecordReader, Map> { - private static final Logger logger = LoggerFactory.getLogger(ColumnFamilyRecordReader.class); + private static final Logger logger = LoggerFactory.getLogger(CqlPagingRecordReader.class); public static final int DEFAULT_CQL_PAGE_LIMIT = 1000; // TODO: find the number large enough but not OOM @@ -96,7 +96,7 @@ public class ColumnFamilyRecordReader extends RecordReader keyValidator; - public ColumnFamilyRecordReader() + public CqlPagingRecordReader() { super(); } @@ -111,12 +111,12 @@ public class ColumnFamilyRecordReader extends RecordReader * - * @see ColumnFamilyOutputFormat + * @see CqlOutputFormat */ -final class ColumnFamilyRecordWriter extends AbstractColumnFamilyRecordWriter, List> +final class CqlRecordWriter extends AbstractColumnFamilyRecordWriter, List> { - private static final Logger logger = LoggerFactory.getLogger(ColumnFamilyRecordWriter.class); + private static final Logger logger = LoggerFactory.getLogger(CqlRecordWriter.class); // handles for clients for each range running in the threadpool private final Map clients; @@ -84,29 +84,29 @@ final class ColumnFamilyRecordWriter extends AbstractColumnFamilyRecordWriter(); - cql = CQLConfigHelper.getOutputCql(conf); + cql = CqlConfigHelper.getOutputCql(conf); try { String host = getAnyHost(); int port = ConfigHelper.getOutputRpcPort(conf); - Cassandra.Client client = ColumnFamilyOutputFormat.createAuthenticatedClient(host, port, conf); + Cassandra.Client client = CqlOutputFormat.createAuthenticatedClient(host, port, conf); retrievePartitionKeyValidator(client); if (client != null) @@ -250,7 +250,7 @@ final class ColumnFamilyRecordWriter extends AbstractColumnFamilyRecordWriter