请求:提供Scala/Java实现Google PubSub到BigQuery流式处理示例
Absolutely, let's build a clean, production-ready streaming pipeline for Pub/Sub to BigQuery using Scala (our priority) and Java. I'll use Apache Beam with Google Cloud Dataflow here—it handles scaling, fault tolerance, and managed streaming out of the box, so you can focus on your core data transformation logic.
Scala Implementation
Step 1: Add Dependencies (sbt)
First, add these to your build.sbt to pull in Beam and Google Cloud libraries:
val beamVersion = "2.54.0" val googleCloudVersion = "2.28.0" libraryDependencies ++= Seq( "org.apache.beam" % "beam-sdks-java-core" % beamVersion, "org.apache.beam" % "beam-runners-google-cloud-dataflow-java" % beamVersion, "org.apache.beam" % "beam-sdks-java-io-google-cloud-platform" % beamVersion, "com.google.cloud" % "google-cloud-bigquery" % googleCloudVersion, "org.scala-lang.modules" %% "scala-java8-compat" % "1.0.2" // For smooth Java 8+ interoperability )
Step 2: Pipeline Code
This example reads from a Pub/Sub subscription, transforms raw messages to your unified schema, and streams them to BigQuery:
import org.apache.beam.sdk.Pipeline import org.apache.beam.sdk.io.gcp.pubsub.PubsubIO import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO import org.apache.beam.sdk.options.{PipelineOptions, PipelineOptionsFactory} import org.apache.beam.sdk.transforms.{DoFn, ParDo} import com.google.api.services.bigquery.model.{TableSchema, TableFieldSchema, TableRow} import java.util.Arrays // Define your raw Pub/Sub message structure (match your actual message format) case class RawPubSubData(userId: String, eventType: String, timestamp: Long, rawPayload: String) // Your unified BigQuery-compatible structure case class UnifiedData( userId: String, eventCategory: String, eventTimestamp: String, processedPayload: String ) object PubSubToBigQueryPipeline { def main(args: Array[String]): Unit = { // Initialize pipeline options (pass args like project ID, region via command line) val options = PipelineOptionsFactory.fromArgs(args).create().as(classOf[PipelineOptions]) options.setRunner(classOf[org.apache.beam.runners.dataflow.DataflowRunner]) val pipeline = Pipeline.create(options) // 1. Read messages from Pub/Sub Subscription val pubsubMessages = pipeline.apply( "Read from Pub/Sub", PubsubIO.readStrings().fromSubscription("projects/YOUR_PROJECT/subscriptions/YOUR_SUBSCRIPTION") ) // 2. Transform raw messages to your unified schema val unifiedData = pubsubMessages.apply( "Transform to Unified Structure", ParDo.of(new DoFn[String, UnifiedData]() { @ProcessElement def processElement(ctx: ProcessContext): Unit = { val rawMessage = ctx.element() // TODO: Replace with your actual parsing logic (use Circe/Play JSON for JSON messages) val rawData = parseRawMessage(rawMessage) val unified = UnifiedData( userId = rawData.userId, eventCategory = mapEventType(rawData.eventType), // Map event types to categories eventTimestamp = convertToIsoTimestamp(rawData.timestamp), // Convert to BigQuery-compatible timestamp processedPayload = cleanPayload(rawData.rawPayload) // Clean/transform payload ) ctx.output(unified) } }) ) // 3. Stream transformed data to BigQuery unifiedData.apply( "Write to BigQuery", BigQueryIO.writeTableRows() .to("YOUR_PROJECT:YOUR_DATASET.YOUR_TABLE") .withSchema(getBigQuerySchema()) // Define your table's schema .withFormatFunction(data => { // Convert Scala case class to BigQuery TableRow val row = new TableRow() row.set("userId", data.userId) row.set("eventCategory", data.eventCategory) row.set("eventTimestamp", data.eventTimestamp) row.set("processedPayload", data.processedPayload) row }) .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND) // Append to existing data .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED) // Create table if missing ) // Run the pipeline (blocks until completion or cancellation) pipeline.run().waitUntilFinish() } // --- Helper Functions (Implement these based on your needs) --- private def parseRawMessage(raw: String): RawPubSubData = { // Example using Circe for JSON parsing: // import io.circe.generic.auto._ // import io.circe.parser.decode // decode[RawPubSubData](raw).getOrElse(throw new IllegalArgumentException("Invalid message format")) RawPubSubData("test-user", "click", System.currentTimeMillis(), "{\"action\":\"button_click\"}") } private def mapEventType(eventType: String): String = { eventType match { case "click" | "view" => "engagement" case "purchase" | "checkout" => "conversion" case _ => "other" } } private def convertToIsoTimestamp(timestamp: Long): String = { java.time.Instant.ofEpochMilli(timestamp).toString() } private def cleanPayload(payload: String): String = { payload.trim.replaceAll("\\s+", " ") // Example: Remove extra whitespace // Add logic to redact sensitive data, normalize format, etc. } private def getBigQuerySchema(): TableSchema = { val fields = Arrays.asList( new TableFieldSchema().setName("userId").setType("STRING").setMode("REQUIRED"), new TableFieldSchema().setName("eventCategory").setType("STRING").setMode("REQUIRED"), new TableFieldSchema().setName("eventTimestamp").setType("TIMESTAMP").setMode("REQUIRED"), new TableFieldSchema().setName("processedPayload").setType("STRING").setMode("NULLABLE") ) new TableSchema().setFields(fields) } }
Java Implementation
If you need a Java version, here's a concise equivalent:
Step 1: Add Dependencies (Maven)
<dependencies> <dependency> <groupId>org.apache.beam</groupId> <artifactId>beam-sdks-java-core</artifactId> <version>2.54.0</version> </dependency> <dependency> <groupId>org.apache.beam</groupId> <artifactId>beam-runners-google-cloud-dataflow-java</artifactId> <version>2.54.0</version> </dependency> <dependency> <groupId>org.apache.beam</groupId> <artifactId>beam-sdks-java-io-google-cloud-platform</artifactId> <version>2.54.0</version> </dependency> <dependency> <groupId>com.google.cloud</groupId> <artifactId>google-cloud-bigquery</artifactId> <version>2.28.0</version> </dependency> <dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> <version>2.15.2</version> <!-- For JSON parsing --> </dependency> </dependencies>
Step 2: Pipeline Code
import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.io.gcp.pubsub.PubsubIO; import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.ParDo; import com.google.api.services.bigquery.model.TableRow; import com.google.api.services.bigquery.model.TableSchema; import com.google.api.services.bigquery.model.TableFieldSchema; import com.fasterxml.jackson.databind.ObjectMapper; import java.util.Arrays; // Unified data structure for BigQuery class UnifiedData { private final String userId; private final String eventCategory; private final String eventTimestamp; private final String processedPayload; public UnifiedData(String userId, String eventCategory, String eventTimestamp, String processedPayload) { this.userId = userId; this.eventCategory = eventCategory; this.eventTimestamp = eventTimestamp; this.processedPayload = processedPayload; } // Getters public String getUserId() { return userId; } public String getEventCategory() { return eventCategory; } public String getEventTimestamp() { return eventTimestamp; } public String getProcessedPayload() { return processedPayload; } } // Raw Pub/Sub message structure class RawPubSubData { public String userId; public String eventType; public long timestamp; public String rawPayload; } public class PubSubToBigQueryJavaPipeline { private static final ObjectMapper MAPPER = new ObjectMapper(); public static void main(String[] args) { PipelineOptions options = PipelineOptionsFactory.fromArgs(args).create(); options.setRunner(org.apache.beam.runners.dataflow.DataflowRunner.class); Pipeline pipeline = Pipeline.create(options); // 1. Read from Pub/Sub pipeline.apply("Read PubSub Messages", PubsubIO.readStrings().fromSubscription("projects/YOUR_PROJECT/subscriptions/YOUR_SUBSCRIPTION")) // 2. Transform data .apply("Transform to Unified Schema", ParDo.of(new DoFn<String, UnifiedData>() { @ProcessElement public void processElement(ProcessContext ctx) throws Exception { RawPubSubData rawData = MAPPER.readValue(ctx.element(), RawPubSubData.class); UnifiedData unified = new UnifiedData( rawData.userId, mapEventType(rawData.eventType), java.time.Instant.ofEpochMilli(rawData.timestamp).toString(), cleanPayload(rawData.rawPayload) ); ctx.output(unified); } })) // 3. Write to BigQuery .apply("Stream to BigQuery", BigQueryIO.writeTableRows() .to("YOUR_PROJECT:YOUR_DATASET.YOUR_TABLE") .withSchema(getBigQuerySchema()) .withFormatFunction(unified -> { TableRow row = new TableRow(); row.set("userId", unified.getUserId()); row.set("eventCategory", unified.getEventCategory()); row.set("eventTimestamp", unified.getEventTimestamp()); row.set("processedPayload", unified.getProcessedPayload()); return row; }) .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND) .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED)); pipeline.run().waitUntilFinish(); } // --- Helpers --- private static String mapEventType(String eventType) { return switch (eventType) { case "click", "view" -> "engagement"; case "purchase", "checkout" -> "conversion"; default -> "other"; }; } private static String cleanPayload(String payload) { return payload.trim().replaceAll("\\s+", " "); } private static TableSchema getBigQuerySchema() { return new TableSchema().setFields(Arrays.asList( new TableFieldSchema().setName("userId").setType("STRING").setMode("REQUIRED"), new TableFieldSchema().setName("eventCategory").setType("STRING").setMode("REQUIRED"), new TableFieldSchema().setName("eventTimestamp").setType("TIMESTAMP").setMode("REQUIRED"), new TableFieldSchema().setName("processedPayload").setType("STRING").setMode("NULLABLE") )); } }
Key Notes
- Replace Placeholders: Swap out
YOUR_PROJECT,YOUR_SUBSCRIPTION,YOUR_DATASET, andYOUR_TABLEwith your actual resources. - Parsing Logic: Update the
parseRawMessage(Scala) or JSON parsing (Java) to match your Pub/Sub message format (JSON, Protobuf, etc.). - Permissions: Ensure the service account running the pipeline has:
- Pub/Sub Subscription Reader access
- BigQuery Data Editor access to your dataset
- Local Testing: To test locally, switch the runner to
DirectRunnerinstead ofDataflowRunner.
内容的提问来源于stack exchange,提问作者ryekos
相关产品推荐
相关产品推荐

