Flink连接Cassandra实现事件富化:求DataStream API的Scala示例
使用Flink DataStream API从Cassandra补全事件流的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
相关产品推荐
相关产品推荐

