Apache NiFi使用Groovy脚本转换CSV为XLS文件失败求助
Apache NiFi CSV转XLS脚本报错问题解决
问题背景
使用Apache NiFi提取SQL数据并导出为XLS文件,流程为:ExecuteQuery(Avro格式)→ CSV RecordWriter(转CSV)→ Groovy脚本(CSV转XLS),但两次脚本均报错。
第一次尝试:脚本与错误
脚本内容
// import org.apache.commons.io.IOUtils import java.nio.charset.* // import java.text.SimpleDateFormat import java.io.* import org.apache.poi.ss.usermodel.* import org.apache.poi.hssf.usermodel.* import org.apache.poi.xssf.usermodel.* import org.apache.poi.ss.util.* import org.apache.poi.ss.usermodel.* import org.apache.poi.hssf.extractor.* def flowFile = session.get() if(!flowFile) return flowFile = session.write(flowFile, {inputStream, outputStream -> try { inputStream.writeTo(outputStream) } catch(e) { log.error("Error during processing", e) session.transfer(flowFile, REL_FAILURE) } } as StreamCallback) session.transfer(flowFile, REL_SUCCESS)
错误说明
脚本仅将CSV输入流直接透传到输出流,未做任何格式转换。若后续修改文件扩展名为.xls,Excel会因文件实际为CSV格式而报错:无法打开文件,因为文件格式或文件扩展名无效。
第二次尝试:脚本与错误
脚本内容
@Grapes(@Grab(group='org.apache.poi', module='poi-ooxml', version='3.9')) import com.opencsv.CSVReader @Grapes(@Grab(group='com.opencsv', module='opencsv', version='4.2')) import org.apache.poi.ss.usermodel.* import org.apache.poi.xssf.streaming.* import org.apache.poi.hssf.usermodel.* import org.apache.poi.xssf.usermodel.* import org.apache.poi.ss.util.* import org.apache.poi.ss.usermodel.* import org.apache.poi.hssf.extractor.* import java.nio.charset.* import java.io.* import org.apache.commons.io.IOUtils def flowFile = session.get() def date = new Date() if(!flowFile) return flowFile = session.write(flowFile, {inputStream, outputStream -> SXSSFSheet sheet1 = null; CSVReader reader = null; Workbook wb = null; String generatedXlsFilePath = "/home/"; FileOutputStream fileOutputStream = null; def filename = flowFile.getAttribute('filename') def path = flowFile.getAttribute('path') def nextLine = '' reader = new CSVReader(new FileReader(path+filename), ','); wb = new SXSSFWorkbook(inputStream); sheet1 = (SXSSFSheet) wb.createSheet('Sheet'); def rowNum = 0; while((nextLine = reader.readNext()) != null) { Row currentRow = sheet1.createRow(rowNum++); for(int i=0; i < nextLine.length; i++) { if(NumberUtils.isDigits(nextLine[i])) { currentRow.createCell(i).setCellValue(Integer.parseInt(nextLine[i])); } else if (NumberUtils.isNumber(nextLine[i])) { currentRow.createCell(i).setCellValue(Double.parseDouble(nextLine[i])); } else { currentRow.createCell(i).setCellValue(nextLine[i]); } } } generatedXlsFilePath = generatedXlsFilePath + 'SAISIE_MAGASING.XLS' outputStream = new FileOutputStream(generatedXlsFilePath.trim()); wb.write(outputStream); wb.close(); outputStream.close(); reader.close(); inputStream.close(); } as StreamCallback) flowFile = session.putAttribute(flowFile, 'filename', filename)
错误说明
- FileNotFoundException:尝试用
FileReader(path+filename)读取本地文件,但NiFi中flowFile的path和filename是属性,并非本地文件系统路径,无法直接读取。 - SXSSFWorkbook构造错误:
new SXSSFWorkbook(inputStream)是错误用法,该构造器用于读取现有Excel流,而非创建新工作簿处理CSV。 - NumberUtils未导入:代码中使用
NumberUtils但未导入对应包,导致编译错误。 - 违反NiFi流处理机制:重新赋值
outputStream为本地文件输出流,未将Excel内容写入NiFi提供的outputStream,导致flowFile无有效内容。
正确解决方案脚本
以下脚本修复了上述所有问题,实现CSV到XLS的正确转换:
// 导入依赖包,若NiFi环境已包含这些包可省略@Grab @Grapes([ @Grab(group='org.apache.poi', module='poi-ooxml', version='4.1.2'), @Grab(group='com.opencsv', module='opencsv', version='5.6'), @Grab(group='org.apache.commons', module='commons-lang3', version='3.12.0') ]) import com.opencsv.CSVReader import org.apache.poi.xssf.streaming.SXSSFWorkbook import org.apache.poi.ss.usermodel.Row import org.apache.poi.ss.usermodel.Cell import org.apache.commons.lang3.math.NumberUtils import java.io.InputStreamReader import java.nio.charset.StandardCharsets def flowFile = session.get() if (!flowFile) return try { flowFile = session.write(flowFile, { inputStream, outputStream -> // 从NiFi输入流读取CSV,指定UTF-8编码 CSVReader reader = new CSVReader(new InputStreamReader(inputStream, StandardCharsets.UTF_8)) SXSSFWorkbook workbook = new SXSSFWorkbook() def sheet = workbook.createSheet("Data") String[] nextLine int rowNum = 0 while ((nextLine = reader.readNext()) != null) { Row row = sheet.createRow(rowNum++) for (int i = 0; i < nextLine.length; i++) { Cell cell = row.createCell(i) String value = nextLine[i] // 尝试识别数字类型 if (NumberUtils.isDigits(value)) { cell.setCellValue(Integer.parseInt(value)) } else if (NumberUtils.isParsable(value)) { cell.setCellValue(Double.parseDouble(value)) } else { cell.setCellValue(value) } } } // 将Excel写入NiFi输出流 workbook.write(outputStream) // 关闭资源 reader.close() workbook.dispose() // SXSSFWorkbook需要释放临时文件 } as StreamCallback) // 修改文件名后缀为.xlsx def originalFilename = flowFile.getAttribute('filename') def newFilename = originalFilename.replaceAll(/\.csv$/, '.xlsx') flowFile = session.putAttribute(flowFile, 'filename', newFilename) session.transfer(flowFile, REL_SUCCESS) } catch (Exception e) { log.error("CSV转XLS失败", e) session.transfer(flowFile, REL_FAILURE) }
关键修复点
- 读取CSV流:通过
InputStreamReader直接读取NiFi提供的inputStream,无需访问本地文件。 - 正确创建Excel工作簿:使用
new SXSSFWorkbook()创建新的流式工作簿,适合处理大数据量。 - 导入NumberUtils:添加
org.apache.commons.lang3.math.NumberUtils的导入,解决类型判断问题。 - 写入NiFi输出流:直接将Excel内容写入参数提供的
outputStream,符合NiFi的流处理机制。 - 修改文件名:自动将
.csv后缀改为.xlsx,确保文件格式与扩展名匹配。 - 异常处理:统一捕获异常,转移到失败关系,并记录错误日志。
内容的提问来源于stack exchange,提问作者O. Sam
相关产品推荐
相关产品推荐

