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

无法通过nifi-flink连接器连接Apache Flink与NiFi源,本地文件转Flink传输问题排查

Hi there! Let's break down your issues step by step since you're new to NiFi and Flink—setting up Site-to-Site (S2S) connections can have a few hidden gotchas, but we'll get this sorted out.

First: Verify Your NiFi Pipeline Configuration

Let's rule out NiFi-side issues first, since that's often where data gets stuck:

  • Check the GetFile Processor:
    • Confirm the Input Directory points to your local folder with files, and the folder actually contains files you want to transfer.
    • Make sure the processor is in a Running state (green icon in NiFi UI), and check the processor's metrics: if FlowFiles In and FlowFiles Out are both 0, GetFile isn't picking up any files (double-check permissions on the input folder too!).
  • Check the "Data For Flink" Output Port:
    • Ensure the port is Running (right-click it → Start if it's grayed out).
    • Critical: Did you connect GetFile's "success" relationship to this output port? It's easy to forget this step—without the connection, data never reaches the port for Flink to pull.
    • Right-click the port → List Queue to see if there are any pending FlowFiles. If there are, NiFi is holding data ready for Flink, so the problem is on the Flink side.

Next, let's validate your Flink setup:

  • Match Connector & Flink Versions: This is a super common pitfall! Make sure your flink-connector-nifi_2.11 version exactly matches your Flink runtime version (e.g., if you're using Flink 1.14.0, use connector version 1.14.0). Mismatched versions can cause silent failures with no error logs.
  • Fix the NiFi URL: Your code uses http://localhost:8080/nifi, but the S2S API endpoint is actually http://localhost:8080/nifi-api (the /nifi path is for the NiFi web UI, not programmatic connections). Update this URL in your SiteToSiteClient.Builder—this is likely the main reason you're not seeing any output!
  • Enable Debug Logs: To see what's happening under the hood, add debug logging to your code before creating the client:
    // Add this at the start of your main method
    org.apache.log4j.BasicConfigurator.configure();
    org.apache.log4j.Logger.getRootLogger().setLevel(org.apache.log4j.Level.DEBUG);
    
    This will show you connection attempts, errors, or port discovery issues that aren't showing up in standard logs.
Site-to-Site Security & Certificate Issues

If you want to use a secure NiFi instance (HTTPS), let's tackle that certificate error:

  • First, Confirm Non-Secure Setup Works: Before diving into security, make sure your non-secure HTTP connection works (following the steps above). Once that's working, you can enable security.
  • Fix the "Requires x509 Certificate" Error:
    When using the TLS Toolkit-generated PKCS12 certificate, you need to convert it to a JKS truststore (or use it directly) for Java to recognize it. Use this keytool command:
    keytool -importkeystore -srckeystore your-cert-file.p12 -srcstoretype PKCS12 -destkeystore truststore.jks -deststoretype JKS
    
    You'll be prompted for the PKCS12 password (the TLS Toolkit generates this—make sure you saved it!) and a new password for the JKS truststore.
  • Configure Flink to Use the Truststore: Add these system properties to your Flink client code before initializing the SiteToSiteClient:
    System.setProperty("javax.net.ssl.trustStore", "/path/to/your/truststore.jks");
    System.setProperty("javax.net.ssl.trustStorePassword", "your-jks-password");
    
    Also, update your NiFi URL to use https://localhost:9443/nifi-api (the default secure port) and set transportProtocol(SiteToSiteTransportProtocol.HTTPS) in your client config.
Quick Troubleshooting Checklist

To wrap up, run through this quick list to cover all bases:

  • GetFile is running and has processed FlowFiles
  • GetFile's success relationship is connected to the "Data For Flink" port
  • The output port is running and has data in its queue (if applicable)
  • Flink connector version matches your Flink version
  • NiFi URL uses /nifi-api instead of /nifi
  • Non-secure setup: NiFi's nifi.remote.input.secure is set to false and nifi.remote.input.http.enabled is true (check nifi.properties)
  • Your machine/IntelliJ can reach NiFi's port (8080 or 9443) with no firewall blocks

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 20:52:51