基于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
相关产品推荐
相关产品推荐

