从 Java 调用 map reduce
为此,您可以在您的 java 应用程序中使用 org.apache.hadoop.mapreduce 命名空间中的类(您可以使用旧的 mapred 使用非常相似的方法,只需检查 API 文档):
Job job = Job.getInstance(new Configuration());
// configure job: set input and output types and directories, etc.
job.setJarByClass(MapReduceCassandra.class);
job.submit();
将数据传递给 mapreduce 作业
如果您的行键集非常小,您可以将其序列化为字符串,并将其作为配置参数传递:
job.getConfiguration().set("CassandraRows", getRowsKeysSerialized()); // TODO: implement serializer
//...
job.submit();
在作业中,您将能够通过上下文对象访问参数:
public void map(
IntWritable key, // your key type
Text value, // your value type
Context context
)
{
// ...
String rowsSerialized = context.getConfiguration().get("CassandraRows");
String[] rows = deserializeRows(rowsSerialized); // TODO: implement deserializer
//...
}
但是,如果您的集合可能是无限的,则将其作为参数传递将是一个坏主意。相反,您应该在文件中传递密钥,并利用分布式缓存。
然后,您可以在提交作业之前将此行添加到上面的部分:
job.addCacheFile(new Path(pathToCassandraKeySetFile).toUri());
//...
job.submit();
在作业中,您将能够通过上下文对象访问此文件:
public void map(
IntWritable key, // your key type
Text value, // your value type
Context context
)
{
// ...
URI[] cacheFiles = context.getCacheFiles();
// find, open and read your file here
// ...
}
注意:所有这些都是针对新 API (org.apache.hadoop.mapreduce)。如果您使用org.apache.hadoop.mapred,则该方法非常相似,但在不同对象上调用了一些相关方法。