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

基于WSO2 Integrator的定时文件采集处理并对接远程API技术咨询

Got it, let's break down how to build this workflow entirely within WSO2 Integrator—no external tools required, just native components and embedded scripting. Here's a practical, step-by-step approach tailored to your needs:

Core Workflow Overview

We’ll use WSO2’s built-in connectors, scheduled tasks, and scripting mediators to create a pipeline that:

  1. Runs on a schedule to scan a target directory
  2. Identifies newly added files
  3. Reads, modifies, and sends content to your remote API
  4. Tracks processed files to avoid duplicates

1. Set Up a Scheduled Polling Trigger

First, create a Scheduled Task to initiate the workflow at your desired interval. This replaces external cron jobs entirely.

  • Go to the WSO2 Integrator Management Console → Tasks → Add Task
  • Choose the Message Injector task type
  • Configure the cron expression (e.g., 0 */5 * * * ? for every 5 minutes)
  • Set the target sequence (we’ll build this next) to trigger on each run

2. Scan for New Files with the File Connector

In your target sequence, use the File Connector to list files in the directory, then filter for only newly added ones using embedded scripting.

Step 2.1: List Files in the Target Directory

Add this XML snippet to your sequence to fetch all files in the directory:

<file.listFiles>
    <directory>/path/to/your/target/directory</directory>
    <filePattern>.*\.txt</filePattern> <!-- Adjust to match your file type (e.g., .csv, .json) -->
</file.listFiles>

This stores the list of files (with metadata like lastModified timestamp) in the fileList context property.

Step 2.2: Filter New Files with Scripting

Use a Script Mediator (Groovy is recommended for better performance in WSO2) to compare file timestamps against the last scan time (stored in WSO2 Registry to persist state across restarts):

// Fetch last scan time from Registry (initialize to 0 if it doesn't exist)
def reg = mc.getRegistry()
def lastScanTime = reg.get('/conf/your-project/lastScanTime') ?: '0'
def currentTime = new Date().getTime().toString()

// Filter files modified after last scan
def fileList = mc.getProperty('fileList')
def newFiles = fileList.findAll { file ->
    file.lastModified > Long.parseLong(lastScanTime)
}

// Store filtered files and update last scan time
mc.setProperty('newFiles', newFiles)
reg.put('/conf/your-project/lastScanTime', currentTime)

3. Read, Modify, and Send File Content

Loop through each new file (use a Foreach Mediator) to process and send the content:

Step 3.1: Read File Content

Inside the foreach loop, use the File Connector to read the file’s content:

<file.readFile>
    <filePath>{$ctx:currentItem/filePath}</filePath>
</file.readFile>

The content will be stored in the fileContent context property.

Step 3.2: Modify Content with Scripting

Add another Script Mediator to adjust the content per your requirements (e.g., replace text, transform JSON/XML, add headers):

def rawContent = mc.getProperty('fileContent')
// Example: Replace placeholder text and format as JSON
def modifiedContent = rawContent.replace('{{CUSTOMER_ID}}', '12345')
def jsonPayload = """{"data": "${modifiedContent}"}"""
mc.setProperty('modifiedPayload', jsonPayload)

Step 3.3: Send to Remote API

Use the HTTP Connector to send the modified payload to your remote endpoint:

<call>
    <endpoint>
        <http method="POST" uri-template="https://your-remote-api.com/submit"/>
    </endpoint>
    <payloadFactory media-type="json">
        <format>$1</format>
        <args>
            <arg evaluator="xml" expression="$ctx:modifiedPayload"/>
        </args>
    </payloadFactory>
</call>

4. Mark Files as Processed (Avoid Duplicates)

After successful API delivery, move the file to a "processed" directory to prevent re-scanning:

<file.moveFile>
    <sourceFile>{$ctx:currentItem/filePath}</sourceFile>
    <destinationDirectory>/path/to/processed/directory</destinationDirectory>
</file.moveFile>

Key Considerations for Production

  • Error Handling: Wrap critical steps in a Try-Catch Mediator to handle failures (e.g., move failed files to an "error" directory and log details with the Log Mediator).
  • Performance: If dealing with large directories, limit the filePattern to specific file types and avoid scanning thousands of files at once.
  • State Management: Using the WSO2 Registry for lastScanTime ensures state persists even if the Integrator restarts. For single-node deployments, you can use a Local Entry instead for faster access.

内容的提问来源于stack exchange,提问作者Pierre Label

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:11:39