如何从Java运行MapReduceIndexerTool并监控状态及设置回调?
针对你提出的关于MapReduceIndexerTool的三个问题,我整理了实用的解决方案,都是生产中常用的场景:
1. 在Java代码中直接运行MapReduceIndexerTool并监控任务状态
MapReduceIndexerTool本身实现了Hadoop的Tool接口,所以可以直接用ToolRunner在Java代码中启动它,同时通过Job对象实时获取任务状态。
示例代码如下:
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.util.ToolRunner; import org.apache.solr.hadoop.MapReduceIndexerTool; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.JobStatus; public class SolrIndexerRunner { public static void main(String[] args) throws Exception { // 初始化Hadoop配置,确保加载集群的core-site.xml、yarn-site.xml等配置 Configuration conf = new Configuration(); // 定义MapReduceIndexerTool的运行参数,和命令行参数一一对应 String[] toolArgs = { "-Dmapred.job.name=Solr_Indexing_Task", "-c", "your_solr_collection", "-i", "/hdfs/input/path", "-o", "/hdfs/output/path", "-m", "2" // 示例:设置map任务数量 }; MapReduceIndexerTool indexerTool = new MapReduceIndexerTool(); indexerTool.setConf(conf); // 启动任务,这里可以选择同步或异步运行 // 同步运行:ToolRunner.run会阻塞直到任务完成,返回退出码 int exitCode = ToolRunner.run(indexerTool, toolArgs); // 快速判断任务结果 if (exitCode == 0) { System.out.println("✅ 索引任务成功完成"); } else { System.out.println("❌ 索引任务失败,退出码:" + exitCode); } // 如果需要实时监控任务运行状态,获取Job实例后循环查询 Job job = indexerTool.getJob(); while (!job.isComplete()) { JobStatus.Status currentState = job.getStatus().getState(); float progress = job.getProgress() * 100; System.out.printf("⏳ 任务当前状态:%s,进度:%.2f%%%n", currentState, progress); Thread.sleep(5000); // 每5秒检查一次 } // 输出最终状态 JobStatus.Status finalState = job.getStatus().getState(); System.out.printf("🏁 任务最终状态:%s%n", finalState); } }
通过这种方式,你可以在Java代码里全程掌控任务的运行状态,包括实时进度和最终结果。
2. 命令行启动任务后,用Java检查任务状态
如果已经通过hadoop jar命令启动了MapReduceIndexerTool,你可以通过Hadoop的JobClient或者YARN REST API来查询任务状态。
方法一:用JobClient查询
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.mapred.JobClient; import org.apache.hadoop.mapred.JobID; import org.apache.hadoop.mapred.JobStatus; public class JobStatusChecker { public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); JobClient jobClient = new JobClient(conf); // 替换成你从命令行控制台或YARN UI拿到的JobID JobID jobId = JobID.forName("job_1690000000000_0001"); // 循环监控直到任务完成 while (true) { JobStatus status = jobClient.getJobStatus(jobId); JobStatus.State state = status.getState(); float progress = status.getProgress() * 100; System.out.printf("任务ID:%s,状态:%s,进度:%.2f%%%n", jobId, state, progress); if (state == JobStatus.State.SUCCEEDED || state == JobStatus.State.FAILED || state == JobStatus.State.KILLED) { break; } Thread.sleep(5000); } System.out.println("任务已结束"); } }
方法二:调用YARN REST API
如果你的集群启用了YARN的REST服务,可以直接通过HTTP请求查询应用状态(MapReduce任务在YARN中属于Application):
import java.net.URI; import java.net.http.HttpClient; import java.net.http.HttpRequest; import java.net.http.HttpResponse; public class YarnApiChecker { public static void main(String[] args) throws Exception { String yarnResourceManagerUrl = "http://your-rm-host:8088/ws/v1/cluster/apps"; String jobId = "job_1690000000000_0001"; HttpClient client = HttpClient.newHttpClient(); HttpRequest request = HttpRequest.newBuilder() .uri(URI.create(yarnResourceManagerUrl + "?applicationId=" + jobId)) .GET() .build(); HttpResponse<String> response = client.send(request, HttpResponse.BodyHandlers.ofString()); System.out.println("YARN API返回结果:" + response.body()); // 解析JSON结果获取状态,比如用Jackson或Gson库 } }
3. 让MapReduce任务完成时发送回调/ Webhook
有几种简单可靠的方式实现这个需求:
方式一:任务完成后主动触发(最直接)
不管是Java直接运行任务还是命令行启动后用Java监控,当检测到任务状态变为SUCCEEDED/FAILED时,直接调用你的Webhook接口:
import java.net.URI; import java.net.http.HttpClient; import java.net.http.HttpRequest; import java.net.http.HttpResponse; // 在任务完成的判断逻辑后添加: if (finalState == JobStatus.Status.SUCCEEDED) { String webhookUrl = "https://your-webhook-endpoint.com/notify"; String payload = String.format("{\"jobId\":\"%s\",\"status\":\"SUCCEEDED\"}", job.getJobID()); HttpClient client = HttpClient.newHttpClient(); HttpRequest request = HttpRequest.newBuilder() .uri(URI.create(webhookUrl)) .POST(HttpRequest.BodyPublishers.ofString(payload)) .header("Content-Type", "application/json") .build(); HttpResponse<String> response = client.send(request, HttpResponse.BodyHandlers.ofString()); System.out.printf("Webhook发送结果:状态码%d,响应:%s%n", response.statusCode(), response.body()); }
方式二:注册JobListener监听任务完成事件
如果是在Java中直接运行任务,可以注册JobListener,当任务完成时自动触发逻辑:
// 在启动任务前注册监听器 Job job = indexerTool.getJob(); job.addJobListener(new JobListener() { @Override public void jobSubmitted(Job job) { System.out.println("任务已提交"); } @Override public void jobCompleted(Job job, JobStatus jobStatus) { String status = jobStatus.getState().toString(); String jobId = job.getJobID().toString(); System.out.printf("任务%s已完成,状态:%s%n", jobId, status); // 这里调用Webhook或执行回调逻辑 sendWebhook(jobId, status); } }); // 启动任务 job.waitForCompletion(true); // 单独封装Webhook发送方法 private static void sendWebhook(String jobId, String status) throws Exception { // 和方式一的代码类似,实现HTTP请求发送 }
方式三:利用Hadoop的自定义OutputFormat(进阶)
如果需要更底层的控制,可以自定义OutputFormat,在任务完成的commit阶段发送事件,但这种方式相对复杂,适合有特殊需求的场景。
内容的提问来源于stack exchange,提问作者Cosmin Ioniță
相关产品推荐
相关产品推荐

