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

Apache NiFi签名验证处理器咨询及自研Java签名验证应用说明

Apache NiFi签名验证处理器技术问题分析与代码参考

最近我在折腾Apache NiFi的签名验证处理器开发,同时自己写了个Java版的签名验证应用,现在想把这套逻辑迁移到NiFi里,或者解决NiFi中签名验证的相关技术痛点。先把我的Java代码贴出来(部分实现),希望能得到针对性的技术建议:

package read_key_pck; 

import static java.nio.charset.StandardCharsets.UTF_8; 
import java.util.Scanner; 
import javax.xml.bind.DatatypeConverter; 
import java.io.FileNotFoundException; 
import java.io.IOException; 
import java.security.KeyFactory; 
import java.security.NoSuchAlgorithmException; 
import java.security.NoSuchProviderException; 
import java.security.PrivateKey; 
import java.security.PublicKey; 
import java.security.Sec...

NiFi签名验证处理器开发核心要点

  • 处理器生命周期管理:NiFi处理器必须继承AbstractProcessor,核心逻辑写在onTrigger方法里处理FlowFile。密钥这类资源建议在onScheduled方法中提前初始化,别每次处理都重新加载,能大幅提升性能。
  • FlowFile内容读写:用ProcessSession的read方法读取待验证内容,验证完成后可以把结果(比如验证状态、签名详情)写入FlowFile属性或者直接修改内容。
  • 配置化密钥管理:别把密钥硬编码进处理器,通过PropertyDescriptor定义密钥文件路径这类配置项,让用户在NiFi UI里配置,既灵活又安全。
  • 异常与路由处理:一定要捕获签名验证相关的异常(比如NoSuchAlgorithmException、InvalidKeyException),根据异常把FlowFile路由到FAILURE关系,同时用NiFi的日志组件记录错误详情,方便排查问题。

基于现有Java代码的NiFi处理器改造示例

把你现有Java代码的核心验证逻辑迁移到NiFi处理器的简化示例如下:

import org.apache.nifi.components.PropertyDescriptor;
import org.apache.nifi.flowfile.FlowFile;
import org.apache.nifi.processor.AbstractProcessor;
import org.apache.nifi.processor.ProcessContext;
import org.apache.nifi.processor.ProcessSession;
import org.apache.nifi.processor.Relationship;
import org.apache.nifi.processor.exception.ProcessException;
import org.apache.nifi.processor.io.InputStreamCallback;
import org.apache.nifi.processor.util.StandardValidators;

import java.io.InputStream;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.security.KeyFactory;
import java.security.PublicKey;
import java.security.spec.X509EncodedKeySpec;
import java.util.*;

public class SignatureValidationProcessor extends AbstractProcessor {

    // 定义公钥文件路径配置项
    public static final PropertyDescriptor PUBLIC_KEY_FILE = new PropertyDescriptor.Builder()
            .name("Public Key File Path")
            .description("Path to the public key file used for signature validation")
            .required(true)
            .addValidator(StandardValidators.FILE_EXISTS_VALIDATOR)
            .build();

    // 定义处理器输出关系:验证成功、验证失败
    public static final Relationship REL_SUCCESS = new Relationship.Builder()
            .name("SUCCESS")
            .description("FlowFiles that passed signature validation")
            .build();
    public static final Relationship REL_FAILURE = new Relationship.Builder()
            .name("FAILURE")
            .description("FlowFiles that failed signature validation")
            .build();

    private Set<Relationship> relationships;
    private List<PropertyDescriptor> propertyDescriptors;
    private PublicKey publicKey;

    @Override
    protected void init(final ProcessorInitializationContext context) {
        final List<PropertyDescriptor> descriptors = new ArrayList<>();
        descriptors.add(PUBLIC_KEY_FILE);
        this.propertyDescriptors = Collections.unmodifiableList(descriptors);

        final Set<Relationship> rels = new HashSet<>();
        rels.add(REL_SUCCESS);
        rels.add(REL_FAILURE);
        this.relationships = Collections.unmodifiableSet(rels);
    }

    @Override
    public Set<Relationship> getRelationships() {
        return relationships;
    }

    @Override
    public List<PropertyDescriptor> getSupportedPropertyDescriptors() {
        return propertyDescriptors;
    }

    @Override
    public void onScheduled(final ProcessContext context) throws ProcessException {
        // 调度时加载公钥,避免重复IO操作
        final String publicKeyPath = context.getProperty(PUBLIC_KEY_FILE).getValue();
        try {
            byte[] publicKeyBytes = Files.readAllBytes(Paths.get(publicKeyPath));
            X509EncodedKeySpec spec = new X509EncodedKeySpec(publicKeyBytes);
            KeyFactory kf = KeyFactory.getInstance("RSA"); // 根据你的签名算法调整
            this.publicKey = kf.generatePublic(spec);
        } catch (Exception e) {
            getLogger().error("Failed to load public key from file {}", new Object[]{publicKeyPath}, e);
            throw new ProcessException(e);
        }
    }

    @Override
    public void onTrigger(final ProcessContext context, final ProcessSession session) throws ProcessException {
        FlowFile flowFile = session.get();
        if (flowFile == null) {
            return;
        }

        final Map<String, String> attributes = new HashMap<>();
        boolean validationPassed = false;

        try {
            // 读取FlowFile中的内容并执行验证
            session.read(flowFile, new InputStreamCallback() {
                @Override
                public void process(InputStream in) throws IOException {
                    String content = new String(in.readAllBytes(), StandardCharsets.UTF_8);
                    // 假设内容格式为「原始数据;Base64编码的签名」,可根据实际约定调整
                    String[] parts = content.split(";");
                    String data = parts[0];
                    String signature = parts[1];

                    // 复用你Java代码中的签名验证逻辑
                    byte[] signatureBytes = DatatypeConverter.parseBase64Binary(signature);
                    java.security.Signature sig = java.security.Signature.getInstance("SHA256withRSA"); // 算法要和签名生成端一致
                    sig.initVerify(publicKey);
                    sig.update(data.getBytes(StandardCharsets.UTF_8));
                    validationPassed = sig.verify(signatureBytes);
                }
            });

            // 根据验证结果路由FlowFile
            if (validationPassed) {
                attributes.put("signature.valid", "true");
                flowFile = session.putAllAttributes(flowFile, attributes);
                session.transfer(flowFile, REL_SUCCESS);
                getLogger().info("Signature validation passed for FlowFile {}", new Object[]{flowFile});
            } else {
                attributes.put("signature.valid", "false");
                flowFile = session.putAllAttributes(flowFile, attributes);
                session.transfer(flowFile, REL_FAILURE);
                getLogger().warn("Signature validation failed for FlowFile {}", new Object[]{flowFile});
            }
        } catch (Exception e) {
            getLogger().error("Error validating signature for FlowFile {}", new Object[]{flowFile}, e);
            attributes.put("signature.error", e.getMessage());
            flowFile = session.putAllAttributes(flowFile, attributes);
            session.transfer(flowFile, REL_FAILURE);
        }
    }
}

额外技术建议

  • 算法一致性:一定要保证NiFi处理器使用的签名算法(比如SHA256withRSA)和生成签名的业务端完全一致,否则必然验证失败。
  • 内容格式约定:提前明确FlowFile中原始数据和签名的分隔规则,也可以把分隔符做成可配置项,提升处理器的通用性。
  • 性能优化:如果要处理高并发的FlowFile,确保密钥实例是线程安全的,或者在onScheduled阶段就完成加载,避免重复IO操作。
  • 调试技巧:用NiFi的GenerateFlowFile生成测试数据,配合AttributeViewer查看验证结果属性,再结合日志排查问题会更高效。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:21:34