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

NiFi使用ExecuteScript处理器爬取网站数据写入FlowFile出现0字节问题

修复方案

核心问题定位

你遇到的空文件问题来自两个原因:

  1. httplib的response.read()返回的是Jython原生字节序列,StringUtil.toBytes()仅支持传入字符串类型参数,直接传入字节序列会出现隐式转换失败,最终写入空内容
  2. 未主动设置FlowFile的mime.type属性,导致NiFi默认将其识别为application/octet-stream

修复后完整代码

#import http.client, urllib.request, urllib.parse, urllib.error, base64
import httplib
#import requests
import java.io
from org.apache.commons.io import IOUtils
from java.nio.charset import StandardCharsets
from org.apache.nifi.processor.io import OutputStreamCallback
import json
from datetime import datetime
from java.lang import Object
from jarray import array

class WriteContentCallback(OutputStreamCallback):
    def __init__(self, content):
        self.content_bytes = content

    def process(self, outputStream):
        try:
            # 直接写入字节序列,不需要做字符串转字节处理
            outputStream.write(self.content_bytes)
        except Exception as error:
            objArray = [error]
            javaArray = array(objArray, Object)
            log.error('Error processing content ===> {}', javaArray)
            raise ValueError

flowFile = None
try:
    flowFile = session.create()
    # 替换为你的实际请求域名和路径
    url_domain = "你的请求域名"
    url = "你的请求路径"
    conn = httplib.HTTPSConnection(url_domain)
    conn.request("GET", url)
    response = conn.getresponse()
    
    # 先判断请求状态,非200直接走失败分支
    if response.status != 200:
        log.error("Request failed with status code: {}".format(response.status))
        session.transfer(flowFile, REL_FAILURE)
    else:
        data = response.read()
        # 根据返回内容设置正确的mime类型,比如返回csv就填text/csv,返回json就填application/json
        flowFile = session.putAttribute(flowFile, "mime.type", "text/csv")
        flowFile = session.write(flowFile, WriteContentCallback(data))
        session.transfer(flowFile, REL_SUCCESS)
    conn.close()
except Exception as outermost_error:
    objArray = [outermost_error]
    javaArray = array(objArray, Object)
    log.error('Error processing request ===> {}', javaArray)
    if flowFile is not None:
        session.transfer(flowFile, REL_FAILURE)

额外注意事项

  • 不要在脚本内硬编码请求地址、鉴权信息等敏感内容,可以通过NiFi表达式语言从处理器属性或者上游FlowFile属性中读取
  • 如果返回内容是GBK、GB2312等非UTF-8编码的文本内容,需要先将字节序列转成对应编码的字符串,再转成字节写入避免乱码
  • 如果需要携带请求头、Cookie等信息,可以在conn.request方法的第四个参数传入对应的头信息字典即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 04:45:07