NiFi执行Groovy脚本转CSV为Excel后生成空FlowFile及属性无效问题
Apache NiFi ExecuteGroovyScript CSV转Excel问题修复
问题原因
FlowFile为空:
- 使用SXSSFWorkbook时直接获取底层XSSFWorkbook操作表格,此时SXSSF的行数据仍在缓存/临时文件中,未同步到底层XSSFWorkbook,导致创建的表格无数据,最终写入输出流时内容丢失。
- 未正确关闭
SXSSFWorkbook和CSVReader,导致输出流未完全刷新。
自定义属性不生效:
- 未通过
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
相关产品推荐
相关产品推荐

