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

基于Spring与GCP构建实时追踪系统:Dataflow集成与数据持久化问题

解决方案:用Google Cloud Dataflow处理PubSub数据到BigQuery & Cloud SQL

首先得明确一个核心点:Dataflow是独立的托管流处理作业,不需要嵌入到你的Spring应用中。你的Spring服务负责从TCP端口接收设备数据并推送到PubSub,而Dataflow则作为单独的流处理链路,订阅PubSub的消息后完成数据转换与双存储写入,两者完全解耦,不需要修改你现有Spring代码的逻辑。

下面分步骤帮你落地:

1. 准备Dataflow开发依赖

如果用Maven管理项目,先在pom.xml中添加Dataflow核心依赖(建议用最新稳定版):

<dependencies>
    <!-- Dataflow GCP SDK -->
    <dependency>
        <groupId>org.apache.beam</groupId>
        <artifactId>beam-sdks-java-google-cloud-platform</artifactId>
        <version>2.54.0</version>
    </dependency>
    <!-- Beam核心工具包 -->
    <dependency>
        <groupId>org.apache.beam</groupId>
        <artifactId>beam-sdks-java-core</artifactId>
        <version>2.54.0</version>
    </dependency>
    <!-- JSON序列化工具(如果你的设备数据是JSON格式) -->
    <dependency>
        <groupId>com.google.code.gson</groupId>
        <artifactId>gson</artifactId>
        <version>2.10.1</version>
    </dependency>
</dependencies>

2. 编写Dataflow Pipeline核心代码

创建一个独立的Java类(比如DeviceTrackingPipeline.java),这是Dataflow作业的入口,负责定义从PubSub读数据、转换格式、写入BigQuery和Cloud SQL的完整链路:

import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO;
import org.apache.beam.sdk.io.gcp.pubsub.PubsubIO;
import org.apache.beam.sdk.io.jdbc.JdbcIO;
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 org.apache.beam.sdk.values.TableRow;
import com.google.gson.Gson;
import com.google.api.services.bigquery.model.TableFieldSchema;
import com.google.api.services.bigquery.model.TableSchema;
import java.util.List;

// 定义与设备数据匹配的POJO类(要和Spring端发送的格式一致)
class DeviceTrackingData {
    private String deviceId;
    private double latitude;
    private double longitude;
    private long timestamp;
    
    // 自动生成getter、setter方法
    public String getDeviceId() { return deviceId; }
    public void setDeviceId(String deviceId) { this.deviceId = deviceId; }
    public double getLatitude() { return latitude; }
    public void setLatitude(double latitude) { this.latitude = latitude; }
    public double getLongitude() { return longitude; }
    public void setLongitude(double longitude) { this.longitude = longitude; }
    public long getTimestamp() { return timestamp; }
    public void setTimestamp(long timestamp) { this.timestamp = timestamp; }
}

public class DeviceTrackingPipeline {
    public static void main(String[] args) {
        // 初始化Pipeline配置(可通过命令行参数传入)
        PipelineOptions options = PipelineOptionsFactory.fromArgs(args).create();
        Pipeline pipeline = Pipeline.create(options);

        // 完整数据流链路
        pipeline.apply("从PubSub拉取数据", PubsubIO.readStrings()
                        .fromSubscription("projects/你的项目ID/subscriptions/exampleSubscription"))
                // 步骤1:将PubSub的字符串消息转为POJO
                .apply("转换为设备数据模型", ParDo.of(new DoFn<String, DeviceTrackingData>() {
                    @ProcessElement
                    public void processElement(ProcessContext ctx) {
                        String jsonMsg = ctx.element();
                        DeviceTrackingData data = new Gson().fromJson(jsonMsg, DeviceTrackingData.class);
                        ctx.output(data);
                    }
                }))
                // 步骤2:分支写入BigQuery
                .apply("写入BigQuery", BigQueryIO.writeTableRows()
                        .to("你的项目ID:你的数据集ID.设备追踪表名")
                        .withSchema(buildBigQuerySchema())
                        .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND)
                        .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED)
                        .withFormatFunction(data -> {
                            TableRow row = new TableRow();
                            row.set("deviceId", data.getDeviceId());
                            row.set("latitude", data.getLatitude());
                            row.set("longitude", data.getLongitude());
                            row.set("timestamp", data.getTimestamp());
                            return row;
                        }))
                // 步骤3:分支写入Cloud SQL(以MySQL为例)
                .apply("写入Cloud SQL", JdbcIO.<DeviceTrackingData>write()
                        .withDataSourceConfiguration(JdbcIO.DataSourceConfiguration.create(
                                "com.mysql.cj.jdbc.Driver",
                                "jdbc:mysql://你的CloudSQLIP:3306/你的数据库名")
                                .withUsername("数据库用户名")
                                .withPassword("数据库密码"))
                        .withStatement("INSERT INTO device_tracking (device_id, latitude, longitude, record_time) VALUES (?, ?, ?, ?)")
                        .withPreparedStatementSetter((data, stmt) -> {
                            stmt.setString(1, data.getDeviceId());
                            stmt.setDouble(2, data.getLatitude());
                            stmt.setDouble(3, data.getLongitude());
                            stmt.setLong(4, data.getTimestamp());
                        }));

        // 启动Pipeline作业
        pipeline.run().waitUntilFinish();
    }

    // 定义BigQuery表结构
    private static TableSchema buildBigQuerySchema() {
        return new TableSchema().setFields(List.of(
                new TableFieldSchema().setName("deviceId").setType("STRING").setMode("REQUIRED"),
                new TableFieldSchema().setName("latitude").setType("FLOAT64").setMode("REQUIRED"),
                new TableFieldSchema().setName("longitude").setType("FLOAT64").setMode("REQUIRED"),
                new TableFieldSchema().setName("timestamp").setType("TIMESTAMP").setMode("REQUIRED")
        ));
    }
}

3. 部署与关键注意事项

部署Dataflow作业

通过Maven命令行部署到GCP(替换占位符为你的实际信息):

mvn compile exec:java \
  -Dexec.mainClass="你的包路径.DeviceTrackingPipeline" \
  -Dexec.args="--project=你的项目ID \
               --runner=DataflowRunner \
               --region=你的GCP区域(比如us-central1) \
               --tempLocation=gs://你的GCS存储桶/temp"

核心注意点

  • 数据格式对齐:确保Spring端发送到PubSub的消息格式(比如JSON)和Dataflow的POJO解析逻辑完全匹配,避免转换失败。
  • Cloud SQL安全连接:不要直接用公网IP连接Cloud SQL,建议用Cloud SQL Auth Proxy或GCP私有网络,Dataflow作业可以通过配置网络参数接入私有网络。
  • 权限配置:给Dataflow服务账号分配以下权限:PubSub订阅读取权、BigQuery写入权、Cloud SQL写入权。
  • 现有Spring应用无需修改:你的Spring服务只需要继续负责TCP数据接收和PubSub推送,Dataflow会自动从exampleSubscription拉取数据,两者互不干扰。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:36:22