【问题标题】:Can't understand TwoPhaseCommitSinkFunction lifecycle无法理解 TwoPhaseCommitSinkFunction 生命周期
【发布时间】:2020-05-14 06:43:56
【问题描述】:

我需要一个连接到 Postgres DB 的接收器,所以我开始构建一个自定义的 Flink SinkFunction。由于 FlinkKafkaProducer 实现了 TwoPhaseCommitSinkFunction,所以我决定做同样的事情。正如 O'Reilley 的书使用 Apache Flink 进行流处理所述,您只需要实现抽象方法,启用检查点,就可以开始了。但是当我运行我的代码时真正发生的事情是 commit 方法只被调用一次,并且在 invoke 之前被调用,这是完全出乎意料的,因为如果你的集合准备就绪,你不应该准备好提交 -提交事务为空。最糟糕的是,在提交之后,我的文件中存在的所有事务行都调用了invoke,然后又调用了abort,这更加出乎意料。

当 Sink 初始化时,据我了解应该会发生以下情况:

  1. beginTransaction 被调用并发送一个标识符来调用
  2. invoke 根据收到的标识符将行添加到事务中
  3. 预提交对当前事务数据进行所有最终修改
  4. commit 处理预提交数据的最终事务

所以,我不明白为什么我的程序没有显示这种行为。

这是我的接收器代码:

package PostgresConnector

import java.sql.{BatchUpdateException, DriverManager, PreparedStatement, SQLException, Timestamp}
import java.text.ParseException
import java.util.{Date, Properties, UUID}
import org.apache.flink.api.common.ExecutionConfig
import org.apache.flink.configuration.Configuration
import org.apache.flink.streaming.api.functions.sink.{SinkFunction, TwoPhaseCommitSinkFunction}
import org.apache.flink.streaming.api.scala._
import org.slf4j.{Logger, LoggerFactory}




class PostgreSink(props : Properties, config : ExecutionConfig) extends TwoPhaseCommitSinkFunction[(String,String,String,String),String,String](createTypeInformation[String].createSerializer(config),createTypeInformation[String].createSerializer(config)){
    
    private var transactionMap : Map[String,Array[(String,String,String,String)]] = Map()
    
    private var parsedQuery : PreparedStatement = _
    
    private val insertionString : String = "INSERT INTO mydb (field1,field2,point) values (?,?,point(?,?))"
    
    override def invoke(transaction: String, value: (String,String,String,String), context: SinkFunction.Context[_]): Unit = {
    
        val LOG = LoggerFactory.getLogger(classOf[FlinkCEPClasses.FlinkCEPPipeline])
        
        val res = this.transactionMap.get(transaction)
        
        if(res.isDefined){
    
            var array = res.get
            
            array = array ++ Array(value)
    
            this.transactionMap += (transaction -> array)
            
        }else{
    
            val array = Array(value)
    
            this.transactionMap += (transaction -> array)
            
            
        }
    
        LOG.info("\n\nPassing through invoke\n\n")
        
        ()
        
    }
    
    override def beginTransaction(): String = {
    
        val LOG: Logger = LoggerFactory.getLogger(classOf[FlinkCEPClasses.FlinkCEPPipeline])
        
        val identifier = UUID.randomUUID.toString
    
        LOG.info("\n\nPassing through beginTransaction\n\n")
        
        identifier
        
        
    }
    
    override def preCommit(transaction: String): Unit = {
        
        val LOG = LoggerFactory.getLogger(classOf[FlinkCEPClasses.FlinkCEPPipeline])
    
        try{
        
            val tuple : Option[Array[(String,String,String,String)]]= this.transactionMap.get(transaction)
        
            if(tuple.isDefined){
            
                tuple.get.foreach( (value : (String,String,String,String)) => {
                
                    LOG.info("\n\n"+value.toString()+"\n\n")
                
                    this.parsedQuery.setString(1,value._1)
                    this.parsedQuery.setString(2,value._2)
                    this.parsedQuery.setString(3,value._3)
                    this.parsedQuery.setString(4,value._4)
                    this.parsedQuery.addBatch()
                
                })
                
            }
        
        }catch{
        
            case e : SQLException =>
                LOG.info("\n\nError when adding transaction to batch: SQLException\n\n")
        
            case f : ParseException =>
                LOG.info("\n\nError when adding transaction to batch: ParseException\n\n")
        
            case g : NoSuchElementException =>
                LOG.info("\n\nError when adding transaction to batch: NoSuchElementException\n\n")
        
            case h : Exception =>
                LOG.info("\n\nError when adding transaction to batch: Exception\n\n")
        
        }
        
        this.transactionMap = this.transactionMap.empty
    
        LOG.info("\n\nPassing through preCommit...\n\n")
    }
    
    override def commit(transaction: String): Unit = {
    
        val LOG : Logger = LoggerFactory.getLogger(classOf[FlinkCEPClasses.FlinkCEPPipeline])
        
        if(this.parsedQuery != null) {
            LOG.info("\n\n" + this.parsedQuery.toString+ "\n\n")
        }
        
        try{
            
            this.parsedQuery.executeBatch
            val LOG : Logger = LoggerFactory.getLogger(classOf[FlinkCEPClasses.FlinkCEPPipeline])
            LOG.info("\n\nExecuting batch\n\n")
            
        }catch{
    
            case e : SQLException =>
                val LOG : Logger = LoggerFactory.getLogger(classOf[FlinkCEPClasses.FlinkCEPPipeline])
                LOG.info("\n\n"+"Error : SQLException"+"\n\n")
            
        }
        
        this.transactionMap = this.transactionMap.empty
    
        LOG.info("\n\nPassing through commit...\n\n")
        
    }
    
    override def abort(transaction: String): Unit = {
    
        val LOG : Logger = LoggerFactory.getLogger(classOf[FlinkCEPClasses.FlinkCEPPipeline])
        
        this.transactionMap = this.transactionMap.empty
    
        LOG.info("\n\nPassing through abort...\n\n")
        
    }
    
    override def open(parameters: Configuration): Unit = {
    
        val LOG: Logger = LoggerFactory.getLogger(classOf[FlinkCEPClasses.FlinkCEPPipeline])
        
        val driver = props.getProperty("driver")
        val url = props.getProperty("url")
        val user = props.getProperty("user")
        val password = props.getProperty("password")
        Class.forName(driver)
        val connection = DriverManager.getConnection(url + "?user=" + user + "&password=" + password)
        this.parsedQuery = connection.prepareStatement(insertionString)
    
        LOG.info("\n\nConfiguring BD conection parameters\n\n")
    }
}

这是我的主程序:

package FlinkCEPClasses

import PostgresConnector.PostgreSink
import org.apache.flink.api.java.io.TextInputFormat
import org.apache.flink.api.java.utils.ParameterTool
import org.apache.flink.cep.PatternSelectFunction
import org.apache.flink.cep.pattern.conditions.SimpleCondition
import org.apache.flink.cep.scala.pattern.Pattern
import org.apache.flink.core.fs.{FileSystem, Path}
import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.api.TimeCharacteristic
import org.apache.flink.cep.scala.{CEP, PatternStream}
import org.apache.flink.streaming.api.functions.source.FileProcessingMode
import org.apache.flink.streaming.api.scala.{DataStream, StreamExecutionEnvironment}
import java.util.Properties

import org.apache.flink.api.common.ExecutionConfig
import org.slf4j.{Logger, LoggerFactory}

class FlinkCEPPipeline {

  val LOG: Logger = LoggerFactory.getLogger(classOf[FlinkCEPPipeline])
  LOG.info("\n\nStarting the pipeline...\n\n")
  
  var env : StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment

  env.enableCheckpointing(10)
  env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime)
  env.setParallelism(1)

  //var input : DataStream[String] = env.readFile(new TextInputFormat(new Path("/home/luca/Desktop/lines")),"/home/luca/Desktop/lines",FileProcessingMode.PROCESS_CONTINUOUSLY,1)

  var input : DataStream[String] = env.readTextFile("/home/luca/Desktop/lines").name("Raw stream")
  
  var tupleStream : DataStream[(String,String,String,String)] = input.map(new S2PMapFunction()).name("Tuple Stream")
  
  var properties : Properties = new Properties()
  
  properties.setProperty("driver","org.postgresql.Driver")
  properties.setProperty("url","jdbc:postgresql://localhost:5432/mydb")
  properties.setProperty("user","luca")
  properties.setProperty("password","root")
  
  tupleStream.addSink(new PostgreSink(properties,env.getConfig)).name("Postgres Sink").setParallelism(1)
  tupleStream.writeAsText("/home/luca/Desktop/output",FileSystem.WriteMode.OVERWRITE).name("File Sink").setParallelism(1)

  env.execute()


}

我的S2PMap函数代码:

package FlinkCEPClasses

import org.apache.flink.api.common.functions.MapFunction

case class S2PMapFunction() extends MapFunction[String,(String,String,String,String)] {
    
    override def map(value: String): (String, String, String,String) = {
    
    
            var tuple = value.replaceAllLiterally("(","").replaceAllLiterally(")","").split(',')
    
            (tuple(0),tuple(1),tuple(2),tuple(3))
        
    }
}

我的管道是这样工作的:我从文件中读取行,将它们映射到字符串元组,然后使用元组中的数据将它们保存在 Postgres 数据库中

如果你想模拟数据,只需创建一个包含如下格式的行的文件: (field1,field2,pointx,pointy)

编辑

TwoPhaseCommitSinkFUnction 的方法执行顺序如下:

Starting pipeline...
beginTransaction
preCommit
beginTransaction
commit
invoke
invoke
invoke
invoke
invoke
invoke
invoke
invoke
invoke
invoke
invoke
invoke
invoke
invoke
invoke
invoke
invoke
invoke
abort

【问题讨论】:

  • 编写和维护一个两阶段提交接收器需要做很多工作。为什么不直接使用带有 JDBC 连接器和 postgres 驱动程序的表 API?见ci.apache.org/projects/flink/flink-docs-stable/dev/table/…
  • @DavidAnderson 谢谢你,伙计!花了几个小时尝试设置表 API,但最终将流数据插入到我的数据库中。我要写一个答案来解释我是如何实现它的,因为我在互联网上没有找到任何文章以直截了当的方式解释它。我不得不在 github 上搜索一些更新的 flink 库代码才能找到我想要的。我将对 TwoPhaseCommitSinkFunction 做同样的事情,因为有人可能已经编写了自定义 flink Sink。
  • 很高兴你有一些工作。至于您在两阶段提交接收器中遇到的问题,我不确定发生了什么,但我注意到您的检查点间隔非常短(10 毫秒)。如果将其增加到更大的值(例如 10 秒),您仍然会得到相同的行为吗?
  • 不,行为改变了。实际上,当我将检查点时间设置为 1000 毫秒时,只会执行 beginTransaction。只有在将检查点时间更改为 10 毫秒后,我才能看到正在执行的其他方法的日志记录。我的猜测是管道处理速度太快,在第一个检查点1000ms之前结束执行方式。
  • 你知道哪个事务被中止了吗?

标签: scala intellij-idea stream apache-flink


【解决方案1】:

所以,这里是这个问题的“答案”。需要明确的是:目前,关于TwoPhaseCommitSinkFunction 的问题还没有解决。如果您正在寻找的是关于原始问题的内容,那么您应该寻找另一个答案。如果您不关心将什么用作水槽,那么也许我可以为您提供帮助。

按照@DavidAnderson 的建议,我开始研究Table API 看看它是否可以解决我的问题,即使用Flink 在我的数据库表中插入行。

如您所见,结果非常简单。

OBS:注意您使用的版本。我的 Flink 版本是1.9.0

Source code

package FlinkCEPClasses

import java.sql.Timestamp
import java.util.Properties

import org.apache.flink.api.common.typeinfo.{TypeInformation, Types}
import org.apache.flink.api.java.io.jdbc.JDBCAppendTableSink
import org.apache.flink.streaming.api.TimeCharacteristic
import org.apache.flink.streaming.api.scala.{DataStream, StreamExecutionEnvironment}
import org.apache.flink.table.api.{EnvironmentSettings, Table}
import org.apache.flink.table.api.scala.StreamTableEnvironment
import org.apache.flink.streaming.api.scala._
import org.apache.flink.table.sinks.TableSink
import org.postgresql.Driver

class TableAPIPipeline {

    // --- normal pipeline initialization in this block ---

    var env : StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment

    env.enableCheckpointing(10)
    env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime)
    env.setParallelism(1)


    var input : DataStream[String] = env.readTextFile("/home/luca/Desktop/lines").name("Original stream")

    var tupleStream : DataStream[(String,Timestamp,Double,Double)] = input.map(new S2PlacaMapFunction()).name("Tuple Stream")

    var properties : Properties = new Properties()

    properties.setProperty("driver","org.postgresql.Driver")
    properties.setProperty("url","jdbc:postgresql://localhost:5432/mydb")
    properties.setProperty("user","myuser")
    properties.setProperty("password","mypassword")

    // --- normal pipeline initialization in this block END ---

    // These two lines create what Flink calls StreamTableEnvironment. 
    // It seems pretty similar to a normal stream initialization.
    val settings = EnvironmentSettings.newInstance().useBlinkPlanner().inStreamingMode().build()
    val tableEnv = StreamTableEnvironment.create(env,settings)

    //Since I wanted to sink data into a database, I used JDBC TableSink,
    //because it is very intuitive and is a exact match with my need. You may
    //look for other TableSink classes that fit better in you solution.
    var tableSink : JDBCAppendTableSink = JDBCAppendTableSink.builder()
    .setBatchSize(1)
    .setDBUrl("jdbc:postgresql://localhost:5432/mydb")
    .setDrivername("org.postgresql.Driver")
    .setPassword("mypassword")
    .setUsername("myuser")
    .setQuery("INSERT INTO mytable (data1,data2,data3) VALUES (?,?,point(?,?))")
    .setParameterTypes(Types.STRING,Types.SQL_TIMESTAMP,Types.DOUBLE,Types.DOUBLE)
    .build()

    val fieldNames = Array("data1","data2","data3","data4")
    val fieldTypes = Array[TypeInformation[_]](Types.STRING,Types.SQL_TIMESTAMP,Types.DOUBLE, Types.DOUBLE)



    // This is the crucial part of the code: first, you need to register
    // your table sink, informing the name, the field names, field types and
    // the TableSink object.

    tableEnv.registerTableSink("postgres-table-sink",
        fieldNames,
        fieldTypes,
        tableSink
    )

    // Then, you transform your DataStream into a Table object.
    var table = tableEnv.fromDataStream(tupleStream)

    // Finally, you insert your stream data into the registered sink.
    table.insertInto("postgres-table-sink")




    env.execute()

}

【讨论】:

    【解决方案2】:

    我不是这个主题的专家,但有几个猜测:

    preCommit 在 Flink 开始检查点时调用,commit 在检查点完成时调用。调用这些方法仅仅是因为检查点正在发生,而不管接收器是否接收到任何数据。

    检查点会定期发生,无论是否有任何数据流经您的管道。鉴于您的检查点间隔非常短(10 毫秒),第一个检查点屏障在源设法向其发送任何数据之前到达接收器似乎是合理的。

    您似乎还假设一次只打开一个事务。我不确定这是否得到严格保证,但只要 maxConcurrentCheckpoints 为 1(这是默认值),你应该没问题。

    【讨论】:

    • 是的,我同意你对preCommitcommit 的看法。但这似乎仍然不能证明我得到的行为是合理的,原因如下:如果您假设流的每个事件都需要超过 10 毫秒才能添加到事务中,那么管道行为一定有问题,因为我没有对其进行任何深入的更改,并且因为 Flink 应该处理至少数万个事件,所以 10 毫秒的时间足以完成一个事务。
    • 如果你假设每个事件的添加时间少于 10 毫秒,那么它也没有意义,因为 10 毫秒后已经有一个填充的事务。唯一有意义的场景是流非常短,只有几个事件(比如文本文件中的 10 行),并且管道执行在第一个检查点之前完成,发送存储在交易到垃圾箱。
    • 我用来测试的文件只有18行。由于它是一个短文件,考虑到 Flink 的速度,它会在不到 0.1ms 的时间内处理完,但我不认为是这样,因为我不相信 Flink 会有诸如破坏文件行为的缺陷太短了,看起来不太对劲。我在问题中添加了我的执行日志。
    • 文件的第一个块包含了全部18行,所以操作系统一次读取整个文件,一旦数据到达Flink,我相信它会很快处理它。我只是想打开文件并开始读取它会产生一些开销,而且可能需要超过 10 毫秒。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-01-03
    • 1970-01-01
    • 2012-11-24
    • 2014-12-05
    • 1970-01-01
    相关资源
    最近更新 更多