【问题标题】:How to programmatically write text to Flink socket?如何以编程方式将文本写入 Flink 套接字?
【发布时间】:2020-05-05 20:35:46
【问题描述】:

我的目标是将字符串“SUCCESS”发送到套接字并让 Apache Flink 的 DataStream 获取该字符串并使用单词“SUCCESS”更新本地文本文件我已经设置了一个“Consumer DataStream" 监视从终端发送到本地主机上的端口 9999 的文本。请参阅下面的代码以获取有效的 Consumer DataStream 代码:​​

package p1;

import org.apache.flink.api.common.functions.FilterFunction;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.core.fs.FileSystem.WriteMode;

public class TestClientStreaming {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStream<String> textStream = env.socketTextStream("localhost", 9999);

        SingleOutputStreamOperator<String> filteredStream = textStream.filter(new FilterFunction<String>() {
            public boolean filter(String value) throws Exception {
                return value.equals("SUCCESS");
            }
        });

        filteredStream.writeAsText(<FILE-PATH>, WriteMode.OVERWRITE);

        env.execute("Reading Flink Stream");
    }
}

这对我有用;当我在终端上运行“nc -l 9999”时,只有在终端中键入“SUCCESS”并按 Enter 键时,FILE-PATH 中的“Consumer Text File”才会更新。 到目前为止很棒!

现在,我想避免使用终端,但仍将文本写入此套接字,以便我的“消费者数据流”可以接收它。该文本可以是任何文本,因为我的“Consumer DataStream”将过滤掉它想要的文本,即“SUCCESS”。

我知道如何写入套接字的唯一两种方法是使用终端(如前所述)或使用Producer DataStream,它不断从“Producer Text File" 并调用 DataStream.writeToSocket() 将它读取的文本写入套接字。然后,我的“消费者数据流”会拾取它,并按预期运行。

是否有任何其他选项可以将文本写入套接字,以便 DataStream 可以拾取该文本?我尝试使用套接字库将“SUCCESS”写入 localhost:9999 但无济于事。

这张图片可能有助于可视化我想要解决的问题以及我已经解决的问题:https://ibb.co/41GzNWM

感谢您的宝贵时间!

【问题讨论】:

  • 请更清楚地说明您的用例。根据当前的描述,您似乎把事情复杂化了。
  • 刚刚更新。如果现在有意义,请告诉我
  • 当您说:“现在,我想避免使用终端,但仍将文本写入此套接字,以便我的“消费者数据流”可以接收它。”这个“文本”的来源是什么?它是如何产生的?人类在某处输入此文本还是以编程方式生成?
  • 人工打字。所以基本上,如果我在终端中键入“RandomText”并按 Enter,我的 Consumer DataStream 会读取它并且什么也不做。但是,如果我在终端中键入“成功”并按 Enter,它会写入“消费者文本文件”
  • @chuckskull 我在帖子末尾添加了一张图片,这可能有助于进一步解释我的问题......或者它可能只会让事情变得更加复杂,但我希望不会!

标签: java sockets apache-flink flink-streaming


【解决方案1】:

您需要创建一个等效的 nc -l aka socket server

示例代码如下:

public class SocketWriter {

    public static void main(String[] args) {
        int port = 9999;
        boolean stop = false;
        ServerSocket serverSocket = null;


        try {

            serverSocket = new ServerSocket(port);

            while (!stop) {
                Socket echoSocket = serverSocket.accept();
                PrintWriter out =
                        new PrintWriter(echoSocket.getOutputStream(), true);
                BufferedReader in =
                        new BufferedReader(
                                new InputStreamReader(echoSocket.getInputStream()));

                out.println("Hello Amazing");
                // do whatever logic you want.

                in.close();
                out.close();
                echoSocket.close();

            }
        }
        catch (IOException e) {
            e.printStackTrace();
        }
        finally {
            try {
                serverSocket.close();
            }
            catch (IOException e) {
                e.printStackTrace();
            }
        }

    }
}

【讨论】:

  • 此外,您可能不需要 Flink 来实现此解决方案。一个简单的客户端/服务器套接字模型可以工作。
  • 感谢您的帮助!是否可以使用此连接多个消费者应用程序?到目前为止我无法弄清楚。 @chuckskull
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2011-06-05
  • 1970-01-01
  • 2013-01-25
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多