Scala局部变量作用域调整:如何重用errormsg实现邮件通知
如何在Akka Streams流中重用
errormsg变量实现邮件通知? 嘿,我来帮你搞定这个变量作用域的问题~ 你现在的errormsg是定义在map的匿名函数内部的,作用域只限于那一小块代码,所以没法在外部拿到它做邮件通知。咱们可以通过调整流的处理结构,把errormsg的作用域扩展到能同时支持写入Kudu和发送邮件的逻辑里,这里有两种靠谱的方案:
方案1:分流处理(推荐,符合流式处理的单一职责原则)
首先,我们可以把需要的错误对象和原始消息封装成一个中间类,这样就能在流的多个分支里访问到errormsg:
// 先定义一个case class来封装需要传递的所有数据 case class ErrorWithPayload(error: Errors, rawMsg: String) // 第一步:转换Kafka消息,生成包含Errors对象和errormsg的中间对象 val errorWithPayloadSource: Source[ErrorWithPayload, Consumer.Control] = kafkaMessages.map(msg => { val bytes: Array[Byte] = msg.record.value() val errormsg = bytes.map(_.toChar).mkString val error = new Errors(1235, "filename", "cdr", "cdr_type", 0, errormsg) ErrorWithPayload(error, errormsg) }) // 第二步:定义两个Sink,分别处理Kudu写入和邮件通知 val kuduSink = new ErrorKuduSink(session, table) val emailNotificationSink = Sink.foreach[ErrorWithPayload](item => { // 在这里写你的邮件通知逻辑,直接用item.rawMsg即可 sendErrorNotificationEmail(item.rawMsg) }) // 用alsoTo把流同时发送到两个Sink,实现并行处理 errorWithPayloadSource.alsoTo(emailNotificationSink).to(kuduSink).run()
这种方式的好处是把Kudu写入和邮件通知的逻辑完全解耦,各自负责单一职责,后期维护和扩展都很方便。如果邮件通知是异步操作(比如调用邮件API返回Future),可以改成Sink.foreachAsync来避免阻塞流:
val emailNotificationSink = Sink.foreachAsync[ErrorWithPayload](parallelism = 3)(item => { // 假设sendErrorNotificationEmailAsync返回Future[Unit] sendErrorNotificationEmailAsync(item.rawMsg) })
方案2:在转换逻辑中直接嵌入邮件通知(不推荐,仅作参考)
如果你暂时不想改动流的结构,也可以直接在原来的map函数里调用邮件通知逻辑,但这种方式会把数据转换和业务逻辑混在一起,后期代码会越来越难维护:
val cdrs: Source[Errors, Consumer.Control] = kafkaMessages.map(msg => { val bytes: Array[Byte] = msg.record.value() val errormsg = bytes.map(_.toChar).mkString // 直接在这里调用邮件通知 sendErrorNotificationEmail(errormsg) new Errors(1235, "filename", "cdr", "cdr_type", 0, errormsg) }) cdrs.to(new ErrorKuduSink(session, table)).run()
额外小提示
- 记得处理邮件通知的错误情况,如果邮件发送失败,要不要重试?可以结合Akka Streams的
RetrySink或者自定义容错逻辑来处理。 - 如果你的
Errors类本身已经包含了errormsg字段(看你代码里确实传了进去),其实也可以不用额外封装ErrorWithPayload,直接从Errors对象里取出消息内容来发邮件,这样更简洁:
val cdrs: Source[Errors, Consumer.Control] = kafkaMessages.map(msg => { val bytes: Array[Byte] = msg.record.value() val errormsg = bytes.map(_.toChar).mkString new Errors(1235, "filename", "cdr", "cdr_type", 0, errormsg) }) cdrs.alsoTo(Sink.foreach(error => sendErrorNotificationEmail(error.errormsgField))).to(kuduSink).run() // 这里的error.errormsgField要替换成你Errors类中对应errormsg的字段名
内容的提问来源于stack exchange,提问作者AOm
相关产品推荐
相关产品推荐

