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

Flink连接Cassandra实现事件富化:求DataStream API的Scala示例

我来给你一个完整的可运行示例,实现用Flink处理入站事件流,并异步从Cassandra查询补充信息来 enrich(补全)事件。咱们一步步来:

1. 准备依赖(build.sbt)

首先确保你的项目引入了Flink核心、Cassandra连接器以及异步IO相关的依赖:

val flinkVersion = "1.17.0"
val cassandraVersion = "4.1.3"

libraryDependencies ++= Seq(
  // Flink DataStream 核心
  "org.apache.flink" %% "flink-streaming-scala" % flinkVersion % "provided",
  // Flink Cassandra 连接器
  "org.apache.flink" %% "flink-connector-cassandra" % flinkVersion,
  // Flink 异步IO 支持
  "org.apache.flink" %% "flink-streaming-java" % flinkVersion % "provided",
  // Cassandra Java Driver
  "com.datastax.oss" % "java-driver-core" % cassandraVersion,
  "com.datastax.oss" % "java-driver-query-builder" % cassandraVersion
)

2. 定义数据模型

先定义几个case class来对应不同的数据结构:

// 入站事件(比如从Kafka/其他数据源接收的原始事件)
case class InboundEvent(eventId: String, userId: String, timestamp: Long)

// Cassandra中存储的补充详情
case class EventDetail(userId: String, userNickname: String, userLevel: Int, userRegion: String)

// 补全后的最终事件
case class EnrichedEvent(eventId: String, userId: String, timestamp: Long, nickname: String, level: Int, region: String)

3. Cassandra表结构

提前在Cassandra中创建存储用户详情的表,比如:

CREATE KEYSPACE IF NOT EXISTS event_keyspace WITH replication = {'class': 'SimpleStrategy', 'replication_factor': 1};

CREATE TABLE IF NOT EXISTS event_keyspace.user_details (
    user_id TEXT PRIMARY KEY,
    nickname TEXT,
    level INT,
    region TEXT
);

4. 核心Flink处理逻辑

下面是完整的Scala代码,实现从模拟数据源接收事件,异步查询Cassandra补全信息,最后输出结果:

import org.apache.flink.api.common.eventtime.WatermarkStrategy
import org.apache.flink.configuration.Configuration
import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.api.scala.async.{AsyncFunction, ResultFuture}
import com.datastax.oss.driver.api.core.CqlSession
import com.datastax.oss.driver.api.core.cql.SimpleStatement

import java.util.concurrent.TimeUnit
import scala.concurrent.{ExecutionContext, Future}
import scala.util.{Failure, Success}

object CassandraEventEnrichment {
  def main(args: Array[String]): Unit = {
    // 1. 创建StreamExecutionEnvironment
    val env = StreamExecutionEnvironment.getExecutionEnvironment
    env.setParallelism(2)

    // 2. 模拟入站事件流(实际场景可替换为Kafka/其他数据源)
    val inboundEvents: DataStream[InboundEvent] = env.fromCollection(Seq(
      InboundEvent("event_001", "user_1001", System.currentTimeMillis()),
      InboundEvent("event_002", "user_1002", System.currentTimeMillis()),
      InboundEvent("event_003", "user_1001", System.currentTimeMillis())
    )).assignTimestampsAndWatermarks(
      WatermarkStrategy.forMonotonousTimestamps[InboundEvent]()
        .withTimestampAssigner((event, _) => event.timestamp)
    )

    // 3. 异步查询Cassandra补全事件
    val enrichedEvents: DataStream[EnrichedEvent] = AsyncDataStream.unorderedWait(
      inboundEvents,
      new CassandraAsyncEnricher(),
      5000, // 异步查询超时时间(毫秒)
      TimeUnit.MILLISECONDS,
      100 // 每个并行实例的异步请求并发数
    )

    // 4. 输出结果(实际场景可替换为写入Kafka/数据库等)
    enrichedEvents.print("Enriched Event: ")

    // 执行作业
    env.execute("Cassandra Event Enrichment Job")
  }
}

// 自定义异步查询Cassandra的Function
class CassandraAsyncEnricher extends AsyncFunction[InboundEvent, EnrichedEvent] {
  private var session: CqlSession = _
  // 用于异步执行Cassandra查询的线程池上下文
  private implicit lazy val executionContext: ExecutionContext = ExecutionContext.global

  // 初始化Cassandra连接(在任务启动时执行一次)
  override def open(parameters: Configuration): Unit = {
    session = CqlSession.builder()
      .withKeyspace("event_keyspace")
      .withLocalDatacenter("datacenter1") // 替换为你的Cassandra数据中心名称
      .build()
  }

  // 处理每个入站事件,异步查询Cassandra
  override def asyncInvoke(event: InboundEvent, resultFuture: ResultFuture[EnrichedEvent]): Unit = {
    // 构建CQL查询语句
    val query = SimpleStatement.builder()
      .setQuery("SELECT nickname, level, region FROM user_details WHERE user_id = ?")
      .addPositionalValue(event.userId)
      .build()

    // 异步执行查询
    val queryFuture: Future[EventDetail] = Future {
      val resultSet = session.execute(query)
      val row = resultSet.one()
      if (row != null) {
        EventDetail(
          event.userId,
          row.getString("nickname"),
          row.getInt("level"),
          row.getString("region")
        )
      } else {
        // 如果没有查到用户详情,返回默认值(根据业务需求调整)
        EventDetail(event.userId, "unknown", 0, "unknown")
      }
    }

    // 处理查询结果,关联原事件并返回
    queryFuture.onComplete {
      case Success(detail) =>
        val enriched = EnrichedEvent(
          event.eventId,
          event.userId,
          event.timestamp,
          detail.userNickname,
          detail.userLevel,
          detail.userRegion
        )
        resultFuture.complete(Iterable(enriched))
      case Failure(e) =>
        resultFuture.completeExceptionally(e)
    }
  }

  // 关闭Cassandra连接(任务结束时执行)
  override def close(): Unit = {
    if (session != null) {
      session.close()
    }
  }
}

关键细节说明

  • 异步IO的选择:这里用AsyncDataStream.unorderedWait而不是orderedWait,因为前者性能更高(不需要等待前面的请求完成),如果你的业务对事件顺序没有严格要求,优先用这个;如果必须保持原事件顺序,换成orderedWait即可。
  • Cassandra Session复用:在open方法中初始化Session,避免为每个事件创建新连接,提升性能。
  • 超时与并发控制:设置合理的超时时间和并发请求数,平衡延迟和资源占用。
  • 异常处理:在queryFuture.onComplete中处理查询失败的情况,避免整个作业因为单个查询失败而崩溃。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 15:32:57