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

NiFi执行Groovy脚本转CSV为Excel后生成空FlowFile及属性无效问题

Apache NiFi ExecuteGroovyScript CSV转Excel问题修复

问题原因

  1. FlowFile为空:

    • 使用SXSSFWorkbook时直接获取底层XSSFWorkbook操作表格,此时SXSSF的行数据仍在缓存/临时文件中,未同步到底层XSSFWorkbook,导致创建的表格无数据,最终写入输出流时内容丢失。
    • 未正确关闭SXSSFWorkbook和CSVReader,导致输出流未完全刷新。
  2. 自定义属性不生效:

    • 未通过session.putAttribute()正确修改FlowFile属性(FlowFile是不可变对象,修改后需重新赋值)。

修复方案

方案1:改用XSSFWorkbook(适合数据量不大的场景)

如果处理的CSV数据量不大,直接使用XSSFWorkbook避免流式处理的缓存问题,同时完善资源关闭逻辑:

@Grab("org.apache.poi:poi:5.0.0")
@Grab("org.apache.poi:poi-ooxml:5.0.0")
@Grab("com.opencsv:opencsv:4.6")

import java.io.*;
import org.apache.poi.ss.usermodel.*;
import org.apache.poi.xssf.usermodel.*;
import org.apache.poi.ss.SpreadsheetVersion;
import org.apache.poi.ss.util.AreaReference;
import org.apache.poi.ss.util.CellReference;
import com.opencsv.CSVReader;

def flowFile = session.get()
if(!flowFile) return

// 添加自定义属性
flowFile = session.putAttribute(flowFile, "test", "success")

flowFile = session.write(flowFile, { inputStream, outputStream ->
    CSVReader csvReader = null
    XSSFWorkbook xssfWorkbook = null
    try {
        xssfWorkbook = new XSSFWorkbook();
        csvReader = new CSVReader(new InputStreamReader(inputStream));
        
        XSSFSheet xssfSheet = xssfWorkbook.createSheet("Sheet");

        String[] strHeaders = null;
        String[] dataRow = null;
        int rowNum = 0;
        while ((dataRow = csvReader.readNext()) != null) {
            if (rowNum == 0) strHeaders = dataRow;
            Row currentRow = xssfSheet.createRow(rowNum);
            for (int i = 0; i < dataRow.length; i++) {
                String cellValue = dataRow[i];
                currentRow.createCell(i).setCellValue(cellValue);
            }
            rowNum++;
        }
        
        if (strHeaders != null) {
            int lastRow = rowNum - 1;
            int lastCol = strHeaders.length - 1;
            AreaReference areaReference = new AreaReference(
                new CellReference(0, 0), 
                new CellReference(lastRow, lastCol), 
                SpreadsheetVersion.EXCEL2007
            );

            XSSFTable xssfTable = xssfSheet.createTable(areaReference);
            
            for (int i = 0; i < strHeaders.length; i++) {
                String columnHeader = strHeaders[i];
                if (xssfTable.getCTTable().getTableColumns().getTableColumnList().size() > i) 
                    xssfTable.getCTTable().getTableColumns().getTableColumnList().get(i).setName(columnHeader);
            }
            xssfTable.getCTTable().addNewTableStyleInfo();
            XSSFTableStyleInfo style = (XSSFTableStyleInfo)xssfTable.getStyle();
            style.setName("TableStyleLight9");
            style.setShowColumnStripes(false);
            style.setShowRowStripes(true);
            xssfTable.getCTTable().addNewAutoFilter().setRef(areaReference.formatAsString());
        }

        xssfWorkbook.write(outputStream)
        outputStream.flush()
    } catch (Exception ex) {
        ex.printStackTrace()
        throw ex // 抛出异常让NiFi处理失败流程
    } finally {
        // 关闭资源
        if (csvReader != null) csvReader.close()
        if (xssfWorkbook != null) xssfWorkbook.close()
    }
} as StreamCallback)

session.transfer(flowFile, REL_SUCCESS)

方案2:保留SXSSFWorkbook(适合大数据量场景)

若必须使用流式处理,需避免直接操作底层XSSFWorkbook,改用SXSSF的API,同时确保资源正确关闭:

@Grab("org.apache.poi:poi:5.0.0")
@Grab("org.apache.poi:poi-ooxml:5.0.0")
@Grab("com.opencsv:opencsv:4.6")

import java.io.*;
import org.apache.poi.ss.usermodel.*;
import org.apache.poi.xssf.streaming.*;
import org.apache.poi.ss.SpreadsheetVersion;
import org.apache.poi.ss.util.AreaReference;
import org.apache.poi.ss.util.CellReference;
import com.opencsv.CSVReader;

def flowFile = session.get()
if(!flowFile) return

flowFile = session.putAttribute(flowFile, "test", "success")

flowFile = session.write(flowFile, { inputStream, outputStream ->
    CSVReader csvReader = null
    SXSSFWorkbook sxssfWorkbook = null
    try {
        sxssfWorkbook = new SXSSFWorkbook();
        csvReader = new CSVReader(new InputStreamReader(inputStream));
        
        sxssfWorkbook.setCompressTempFiles(true);
        SXSSFSheet sxssfSheet = sxssfWorkbook.createSheet("Sheet");
        sxssfSheet.setRandomAccessWindowSize(100);

        String[] strHeaders = null;
        String[] dataRow = null;
        int rowNum = 0;
        while ((dataRow = csvReader.readNext()) != null) {
            if (rowNum == 0) strHeaders = dataRow;
            Row currentRow = sxssfSheet.createRow(rowNum);
            for (int i = 0; i < dataRow.length; i++) {
                String cellValue = dataRow[i];
                currentRow.createCell(i).setCellValue(cellValue);
            }
            rowNum++;
        }
        
        // 先将SXSSF数据写入临时输出流,再读取处理表格
        if (strHeaders != null) {
            ByteArrayOutputStream tempOut = new ByteArrayOutputStream()
            sxssfWorkbook.write(tempOut)
            sxssfWorkbook.dispose() // 释放临时文件
            
            // 重新读取字节流为XSSFWorkbook处理表格
            XSSFWorkbook xssfWorkbook = new XSSFWorkbook(new ByteArrayInputStream(tempOut.toByteArray()))
            XSSFSheet xssfSheet = xssfWorkbook.getSheet("Sheet")
            int lastRow = rowNum - 1;
            int lastCol = strHeaders.length - 1;
            AreaReference areaReference = new AreaReference(
                new CellReference(0, 0), 
                new CellReference(lastRow, lastCol), 
                SpreadsheetVersion.EXCEL2007
            );

            XSSFTable xssfTable = xssfSheet.createTable(areaReference);
            
            for (int i = 0; i < strHeaders.length; i++) {
                String columnHeader = strHeaders[i];
                if (xssfTable.getCTTable().getTableColumns().getTableColumnList().size() > i) 
                    xssfTable.getCTTable().getTableColumns().getTableColumnList().get(i).setName(columnHeader);
            }
            xssfTable.getCTTable().addNewTableStyleInfo();
            XSSFTableStyleInfo style = (XSSFTableStyleInfo)xssfTable.getStyle();
            style.setName("TableStyleLight9");
            style.setShowColumnStripes(false);
            style.setShowRowStripes(true);
            xssfTable.getCTTable().addNewAutoFilter().setRef(areaReference.formatAsString());
            
            // 写入最终输出流
            xssfWorkbook.write(outputStream)
            xssfWorkbook.close()
            tempOut.close()
        } else {
            sxssfWorkbook.write(outputStream)
        }
        outputStream.flush()
    } catch (Exception ex) {
        ex.printStackTrace()
        throw ex
    } finally {
        if (csvReader != null) csvReader.close()
        if (sxssfWorkbook != null) {
            sxssfWorkbook.dispose() // 释放SXSSF临时文件
            sxssfWorkbook.close()
        }
    }
} as StreamCallback)

session.transfer(flowFile, REL_SUCCESS)

关键修复点

  • 属性添加:通过flowFile = session.putAttribute(flowFile, "属性名", "属性值")修改属性,重新赋值FlowFile对象。
  • 资源管理:在finally块中关闭所有流和Workbook,确保数据完全写入输出流。
  • SXSSF使用:避免直接操作底层XSSFWorkbook,若需表格功能,先将SXSSF数据写入临时流再转XSSF处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 19:50:33