【问题标题】:Custom Flink sink invoked but with no data调用了自定义 Flink 接收器但没有数据
【发布时间】:2020-08-06 14:56:53
【问题描述】:

我打算实现一个自定义接收器,我在其中创建了一个存根调用函数,该函数仅将接收到的数据记录到任务日志文件(现在),如下所示。

package io.name.package;

import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import org.slf4j.LoggerFactory;
import org.slf4j.Logger;

public class AlertSink extends RichSinkFunction<Alert> {
    Logger LOG = LoggerFactory.getLogger(AlertSink.class);

    @Override
    public void invoke(Alert alert, Context context) throws Exception {
        LOG.info("Invoking sink for alert: ", alert.toString());
    }
}

我已经配置了如下图所示的数据流。

        DataStream<Alert> result = filteredMetrics
            .keyBy(
                new KeySelector<Tuple7<String, String, String, String, String, String, Object>, Tuple3<String, String, String>>() {
                    @Override
                    public  Tuple3<String, String, String> getKey(Tuple7<String, String, String, String, String, String, Object> in) throws Exception {
                        return Tuple3.of(in.f0, in.f1, in.f2);
                    }
            })
            .window(SlidingProcessingTimeWindows.of(Time.seconds(10), Time.seconds(5)))
            .process(new ThresholdEvaluator());

        result.addSink(new AlertSink());

当我检查日志时,我看到接收器已被调用但显示为空字符串。 ThresholdEvaluator 发出警报,但显示非空字符串。

2020-08-05 19:38:16,638 INFO  io.name.package.AlertSink                      - Invoking sink for alert:
2020-08-05 19:38:16,638 INFO  io.name.package.ThresholdEvaluator             - Alert: {"thresholdID":"123123123","grouping":"svc-platform-5445135-production-graph-service-account-toke644zm","period":"5m","isActive":true,"status":"new","firstSeen":1596656296638,"lastSeen":1596656296638,"count":1}
2020-08-05 19:38:16,640 INFO  io.name.package.AlertSink                      - Invoking sink for alert:
2020-08-05 19:38:16,640 INFO  io.name.package.ThresholdEvaluator             - Alert: {"thresholdID":"123123123","grouping":"svc-platform-5445135-staging-input-service-account-token-q57xf","period":"5m","isActive":true,"status":"new","firstSeen":1596656296640,"lastSeen":1596656296640,"count":1}
2020-08-05 19:38:16,643 INFO  io.name.package.AlertSink                      - Invoking sink for alert:
2020-08-05 19:38:16,643 INFO  io.name.package.ThresholdEvaluator             - Alert: {"thresholdID":"123123123","grouping":"svc-platform-5445135-restructure-repo-cmi-service-account-k76cd","period":"5m","isActive":true,"status":"new","firstSeen":1596656296643,"lastSeen":1596656296643,"count":1}
2020-08-05 19:38:16,646 INFO  io.name.package.AlertSink                      - Invoking sink for alert:
2020-08-05 19:38:16,646 INFO  io.name.package.ThresholdEvaluator             - Alert: {"thresholdID":"123123123","grouping":"svc-integrations-14361530-demo-slack-token","period":"5m","isActive":true,"status":"new","firstSeen":1596656296645,"lastSeen":1596656296645,"count":1}

我错过了什么吗?

我还尝试在 ThresholdEvaluator 和 addSink 运算符之间添加 map 函数。 MapFunction 似乎能够很好地接收 Alert 对象,但不能接收 AlertSink。

        result.map(new MapFunction<Alert, Alert>() {
            @Override
            public Alert map(Alert value) {
                LOG.info(value.toString());
                return value;
            }
        }).addSink(new AlertSink());

(使用附加日志更新)

【问题讨论】:

  • 如果您使用result.print() 而不是result.addSink(new AlertSink()),您会得到不为空的结果吗?我想知道问题是否出在其他地方,例如ThresholdEvaluator
  • 我没有尝试 result.print(),但是从 ThresholdEvaluator 我也在发出输出之前记录了输出警报,如下所示。 2020-08-05 19:38:16,646 INFO io.name.package.ThresholdEvaluator - Alert: {"thresholdID":"123123123","grouping":"svc-integrations-14361530-demo-slack-token","period":"5m","isActive":true,"status":"new","firstSeen":1596656296645,"lastSeen":1596656296645,"count":1}
  • 我已使用 ThresholdEvaluator 日志更新了问题,该日志记录了发出的输出(警报)。
  • 我看不出 AlertSink 有什么问题,所以我对 ThresholdEvaluator 保持怀疑。
  • 真的很奇怪。我尝试在 ThresholdEvaluator 和 addSink 运算符之间插入 MapFunction。 MapFunction 能够很好地接收 Alert 对象,但仍然不能接收 AlertSink。

标签: apache-flink flink-streaming


【解决方案1】:

原因是日志输出没有占位符来插值——正确的语法是

LOG.info("Invoking sink for alert: {}", alert.toString());

这个和这个的区别

LOG.info("Invoking sink for alert: " + alert.toString());

是在后一种情况下,无论日志级别如何,每次都会发生字符串连接,并且在第一种情况下,如果只有日志级别至少为INFO,它将插入该值。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-04-27
    • 2021-08-25
    • 2014-04-26
    • 2017-08-19
    • 2020-11-26
    • 1970-01-01
    相关资源
    最近更新 更多