如何让独立运行的Oozie Workflow实现Coordinator的文件校验逻辑?
Got it, let's solve this problem where you want your standalone Oozie Workflow to only execute when the required HDFS input files are present—mirroring the behavior of your Coordinator setup. The mapred.input.dir approach didn’t trigger the pre-execution check you need, so here are three reliable solutions to implement this logic:
1. Pre-Check with a Shell Action (Simplest Approach)
Add a Shell Action as the first step in your Workflow to explicitly validate the HDFS input path. This action will run a shell command to check for files, and only proceed to your main workflow logic if the check passes.
Example Workflow Snippet:
<workflow-app name="check-input-then-run" xmlns="uri:oozie:workflow:0.5"> <start to="check-hdfs-input"/> <!-- Shell Action to validate input files --> <action name="check-hdfs-input"> <shell xmlns="uri:oozie:shell-action:0.3"> <job-tracker>${jobTracker}</job-tracker> <name-node>${nameNode}</name-node> <exec>bash</exec> <argument>-c</argument> <argument> # Check if there are any non-directory files in the input path FILE_COUNT=$(hdfs dfs -ls /path/to/your/input | grep -v "^d" | wc -l) if [ $FILE_COUNT -gt 0 ]; then exit 0 # Success: files exist else exit 1 # Failure: no files found fi </argument> <file>/path/to/hadoop-config/hdfs-site.xml#hdfs-site.xml</file> <file>/path/to/hadoop-config/core-site.xml#core-site.xml</file> </shell> <ok to="main-workflow-action"/> <error to="fail-workflow"/> </action> <!-- Your main workflow logic here (e.g., MapReduce, Spark, etc.) --> <action name="main-workflow-action"> <!-- Replace with your actual action configuration --> <ok to="end"/> <error to="fail-workflow"/> </action> <kill name="fail-workflow"> <message>HDFS input files not found. Workflow terminated.</message> </kill> <end name="end"/> </workflow-app>
- How it works: The shell command counts non-directory entries in the target HDFS path. If the count is greater than 0, the action succeeds and proceeds to your main workflow. If not, it exits with code 1, triggering the
fail-workflowkill node. - Customization: Adjust the check to fit your needs—e.g., check for specific file patterns with
hdfs dfs -test -e /path/to/input/*.csv, or verify minimum file size withhdfs dfs -du /path/to/input | awk '{sum+=$1} END {if(sum > 1048576) exit 0; else exit 1}'.
2. Custom Java Action (For Advanced Validation)
If you need more complex validation (e.g., checking minimum file size, specific file extensions, or matching metadata), build a lightweight Java program that uses HDFS APIs to validate the input path. Oozie will use the program’s exit code to decide whether to proceed.
Example Java Code Snippet:
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import org.apache.hadoop.fs.FileStatus; import java.util.Arrays; public class HdfsInputValidator { public static void main(String[] args) throws Exception { if (args.length != 1) { System.err.println("Usage: HdfsInputValidator <hdfs-input-path>"); System.exit(1); } Configuration conf = new Configuration(); FileSystem fs = FileSystem.get(conf); Path inputPath = new Path(args[0]); // Check if path exists and has at least one valid file if (fs.exists(inputPath)) { FileStatus[] fileStatuses = fs.listStatus(inputPath); boolean hasValidFiles = Arrays.stream(fileStatuses) .filter(status -> !status.isDirectory()) .anyMatch(status -> { // Example: Check for .parquet files larger than 1MB return status.getPath().getName().endsWith(".parquet") && status.getLen() > 1048576; }); if (hasValidFiles) { System.exit(0); // Success: valid files exist } else { System.err.println("No valid .parquet files (>=1MB) found in input path"); System.exit(1); } } else { System.err.println("Input path does not exist"); System.exit(1); } } }
Workflow Configuration for Java Action:
Package the JAR (with dependencies if needed) into HDFS, then add this action to your workflow:
<action name="validate-input-java"> <java xmlns="uri:oozie:java-action:0.2"> <job-tracker>${jobTracker}</job-tracker> <name-node>${nameNode}</name-node> <main-class>com.yourcompany.HdfsInputValidator</main-class> <arg>/path/to/your/input</arg> <archive>hdfs:///path/to/your/validator.jar#validator.jar</archive> <file>/path/to/hadoop-config/hdfs-site.xml#hdfs-site.xml</file> <file>/path/to/hadoop-config/core-site.xml#core-site.xml</file> </java> <ok to="main-workflow-action"/> <error to="fail-workflow"/> </action>
3. Oozie Decision Node with EL Functions (Quick Path Check)
Use Oozie’s built-in Decision Node with EL (Expression Language) functions to check if the HDFS path exists and contains data. Note: This works best if you’re validating a specific file or ensuring a directory isn’t empty.
Example Decision Node Snippet:
<start to="check-input-decision"/> <decision name="check-input-decision"> <switch> <!-- Check if input directory exists and has non-zero total size --> <case to="main-workflow-action">${fs:exists('/path/to/your/input') and fs:fileSize('/path/to/your/input') > 0}</case> <default to="fail-workflow"/> </switch> </decision> <!-- Rest of your workflow (main action, kill, end) -->
- Note:
fs:fileSize()returns the total size of the path. If the directory is empty, this will return 0, so combining it withfs:exists()ensures the directory exists and has actual data.
Why mapred.input.dir Didn’t Work
The mapred.input.dir parameter is passed directly to your MapReduce job, not validated by Oozie upfront. Oozie will start the workflow regardless, and the MapReduce job itself will fail later if no input files are found. This doesn’t block the workflow from starting, which is why you didn’t see the pre-execution check you wanted.
内容的提问来源于stack exchange,提问作者Vishal Singla

