You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让独立运行的Oozie Workflow实现Coordinator的文件校验逻辑?

Standalone Oozie Workflow: Only Run When HDFS Input Files Are Available

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-workflow kill 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 with hdfs 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 with fs: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.22 09:17:43