寻求Hadoop MapReduce任务故障测试示例代码链接
Hey there! As a fellow former MapReduce newbie who’s spent way too much time testing task failures, I’ve got some straightforward, actionable examples to help you experiment with failure scenarios. Below are practical code snippets and tips to simulate different types of task failures:
1. Simulating Map Task Failure
You can intentionally trigger a map task failure by throwing an exception or forcing a process exit in your Mapper class. This lets you see how Hadoop retries failed tasks automatically. Here’s a simple example:
import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import java.io.IOException; import java.util.Random; public class FailingMapper extends Mapper<Object, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text word = new Text(); private Random random = new Random(); @Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { String[] tokens = value.toString().split(" "); for (String token : tokens) { word.set(token); // Simulate random failure: fail 30% of the time if (random.nextFloat() < 0.3) { throw new IOException("Intentional map task failure for testing!"); // Alternatively, force exit to simulate a sudden node crash: // System.exit(1); } context.write(word, one); } } }
By default, Hadoop will retry failed map tasks up to 4 times (controlled by the mapreduce.map.maxattempts configuration). You can adjust this number in your job setup if you want to test fewer/more retries.
2. Simulating Reduce Task Failure
Just like with map tasks, you can inject failures into the Reducer class. Here’s an example where the reducer fails when processing a specific key:
import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; import java.io.IOException; public class FailingReducer extends Reducer<Text, IntWritable, Text, IntWritable> { @Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; for (IntWritable val : values) { sum += val.get(); } // Fail explicitly when processing the key "test" if (key.toString().equals("test")) { throw new RuntimeException("Intentional reduce task failure for key: " + key); } context.write(key, new IntWritable(sum)); } }
When the reducer hits the key "test", it will throw an error and trigger a retry. The default number of reduce task retries is also 4, controlled by mapreduce.reduce.maxattempts.
3. Testing Node-Level Failures (Advanced)
To simulate a worker node crashing mid-task, you can use YARN commands to kill the task’s container:
- First, get your job’s application ID:
yarn application -list - Find the container ID associated with the running task:
yarn container -list -appId <your-application-id> - Kill the container to simulate a node crash:
yarn container -kill <container-id>
Hadoop will detect the lost container and reschedule the task on another available node.
Key Configuration Tweaks for Testing
Adjust these settings in your job configuration to modify failure behavior:
mapreduce.map.maxattempts: Set the maximum number of map task retries (default: 4)mapreduce.reduce.maxattempts: Set the maximum number of reduce task retries (default: 4)mapreduce.task.timeout: Adjust how long Hadoop waits for a task heartbeat before marking it failed (default: 10 minutes)
You can set these directly in your job setup code:
job.setMapMaxAttempts(2); job.setReduceMaxAttempts(2);
内容的提问来源于stack exchange,提问作者Miss.Saung

