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

请求:提供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

  1. Replace Placeholders: Swap out YOUR_PROJECT, YOUR_SUBSCRIPTION, YOUR_DATASET, and YOUR_TABLE with your actual resources.
  2. Parsing Logic: Update the parseRawMessage (Scala) or JSON parsing (Java) to match your Pub/Sub message format (JSON, Protobuf, etc.).
  3. Permissions: Ensure the service account running the pipeline has:
    • Pub/Sub Subscription Reader access
    • BigQuery Data Editor access to your dataset
  4. Local Testing: To test locally, switch the runner to DirectRunner instead of DataflowRunner.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:55:32