Apache Beam 读取 ServiceNow 数据到文本文件:CdapServiceNowToTxt 示例实战指南 【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载导读本文基于 Apache Beam 仓库中的CdapServiceNowToTxt批处理示例讲解如何借助 CDAP 的 ServiceNow 插件通过CdapIO集成从 ServiceNow 实例批量拉取 JSON 格式数据并写入本地.txt文件。读完本文你将掌握如何配置 Gradle 执行任务、如何填写 10 个核心管道参数认证信息、查询模式、值类型、输出路径等、如何理解整条管道的源码级执行链路以及如何切换到不同 Runner 运行该管道。该示例位于 examples/java/cdap/servicenow/src/main/java/org/apache/beam/examples/complete/cdap/servicenow属于 examples 下 CDAP 插件示例集合的一部分cdap 示例总览同目录下还有 Salesforce、Hubspot、Zendesk 等同类示例。一、示例概览从 ServiceNow 到 .txt 的批处理管道CdapServiceNowToTxt是一条纯批处理管道目标非常明确以 JSON 格式从 CDAP ServiceNow 源读取数据把结果记录写入 .txt 文件。ServiceNow 相关参数与输出文件路径均由用户在模板参数中指定。1.1 类定位与入口主类为 CdapServiceNowToTxt.java其main方法负责两件事public static void main(String[] args) { CdapServiceNowOptions options PipelineOptionsFactory.fromArgs(args).withValidation().as(CdapServiceNowOptions.class); // Create the pipeline Pipeline pipeline Pipeline.create(options); run(pipeline, options); }通过PipelineOptionsFactory.fromArgs(args).withValidation()解析命令行参数并开启参数校验——所有标了Validation.Required的参数若缺失会直接报错将参数转换为CdapServiceNowOptions类型后创建Pipeline并调用run(pipeline, options)执行。1.2 三步式执行流程run方法中的注释清晰地概括了管道核心步骤/* * Steps: * 1) Read messages in from Cdap ServiceNow * 2) Extract values only * 3) Write successful records to .txt file */对应到代码实现是一条由三个 Transform 串起来的流水线pipeline .apply(readFromCdapServiceNow, FormatInputTransform.readFromCdapServiceNow(paramsMap)) .setCoder(KvCoder.of( NullableCoder.of(WritableCoder.of(NullWritable.class)), SerializableCoder.of(StructuredRecord.class))) .apply(MapValues.into(TypeDescriptors.strings()) .via(StructuredRecordUtils::structuredRecordToString)) .setCoder(KvCoder.of( NullableCoder.of(WritableCoder.of(NullWritable.class)), StringUtf8Coder.of())) .apply(Values.create()) .apply(writeToTxt, TextIO.write().to(options.getOutputTxtFilePathPrefix()));各环节说明步骤Transform作用1FormatInputTransform.readFromCdapServiceNow(paramsMap)通过CdapIO.read()加载 ServiceNow CDAP 插件读取数据产出KVNullWritable, StructuredRecord2MapValuesStructuredRecordUtils::structuredRecordToString仅提取 Value 部分把StructuredRecord序列化为 JSON 字符串3Values.create()TextIO.write()丢弃 Key把每条 JSON 记录写入以指定前缀命名的.txt文件其中第 2 步用到的structuredRecordToString定义在 StructuredRecordUtils.java内部使用 Gson 将StructuredRecord转为 JSON若记录为null则输出{}保证下游写出时不会出现空记录崩溃。1.3 关键点CDAP 插件驱动的数据源FormatInputTransform.readFromCdapServiceNow定义在 FormatInputTransform.java是数据接入的核心final PluginConfig pluginConfig new ConfigWrapper(ServiceNowSourceConfig.class).withParams(pluginConfigParams).build(); checkStateNotNull(pluginConfig, Plugin config cant be null.); return CdapIO.NullWritable, StructuredRecordread() .withCdapPluginClass(ServiceNowSource.class) .withPluginConfig(pluginConfig) .withKeyClass(NullWritable.class) .withValueClass(StructuredRecord.class);这段代码揭示了底层机制用ConfigWrapperServiceNowSourceConfig把参数 Map 包装成 CDAP 插件配置对象通过CdapIO.read()指定插件类ServiceNowSource来自io.cdap.plugin.servicenow.source包即以 CDAP 生态中的 ServiceNow 批处理 Source 插件作为数据源Key 类型为 Hadoop 的NullWritableValue 类型为 CDAP 的StructuredRecord这与管道中设置的 Coder 一一对应。也就是说该示例本身不直接实现 ServiceNow REST 调用逻辑而是复用 CDAP 数据集成生态中成熟的 ServiceNow 插件这正是 Apache BeamCdapIO桥接层org.apache.beam.sdk.io.cdap的典型用法。二、Gradle 准备创建执行任务在运行示例前需要在项目的build.gradle中声明一个 JavaExec 任务用于动态指定主类与命令行参数task executeCdapServiceNow (type:JavaExec) { mainClass System.getProperty(mainClass) classpath sourceSets.main.runtimeClasspath systemProperties System.getProperties() args System.getProperty(exec.args, ).split() }该任务的核心设计思路配置项含义mainClass System.getProperty(mainClass)从系统属性读取主类全限定名方便在命令行用-DmainClass...动态指定classpath sourceSets.main.runtimeClasspath使用主源码集的运行时 classpath保证 CDAP 插件、Beam SDK 等依赖可用systemProperties System.getProperties()把 JVM 系统属性透传给子进程args System.getProperty(exec.args, ).split()把-Dexec.args传入的参数字符串按空格拆分成参数数组未提供时为空三、运行 CdapServiceNowToTxt 管道3.1 基本运行命令Gradle 任务executeCdapServiceNow支持通过如下命令运行管道gradle clean executeCdapServiceNow -DmainClassorg.apache.beam.examples.complete.cdap.servicenow.CdapServiceNowToTxt \ -Dexec.args--argumentvalue --argumentvalue要点说明clean用于清理上次构建产物避免旧 class 干扰-DmainClass指定入口类org.apache.beam.examples.complete.cdap.servicenow.CdapServiceNowToTxt-Dexec.args内以--keyvalue形式传递所有管道参数参数之间用空格分隔。3.2 完整参数示例执行管道时需按如下格式指定参数--clientIdyour-client-id \ --clientSecretyour-client-secret \ --useryour-user \ --passwordyour-password \ --restApiEndpointyour-endpoint \ --queryModeTable \ --tableNameyour-table \ --valueTypeActual \ --referenceNameyour-reference-name \ --outputTxtFilePathPrefixyour-path-to-output-folder-with-filename-prefix3.3 切换 Runner默认情况下管道使用 DirectRunner 在本地执行。如需更换执行引擎追加--runnerYOUR_SELECTED_RUNNER例如在本地验证时用--runnerDirectRunner默认在分布式环境可换成DataflowRunner、FlinkRunner、SparkRunner等 Apache Beam 支持的 Runner。四、参数详解9 个必填参数 1 个参考名所有参数定义集中在 CdapServiceNowOptions.java该接口继承自 BaseCdapOptions.java后者统一提供了referenceName参数。下面逐一解析4.1 认证类参数4 个参数类型必填说明clientIdString✅ServiceNow 实例的 Client IDOAuth 客户端标识clientSecretString✅ServiceNow 实例的 Client SecretOAuth 客户端密钥userString✅ServiceNow 实例的用户名passwordString✅ServiceNow 实例的密码源码中以Validation.Required标注运行时若缺失会触发校验失败。安全提示明文在命令行中传递密码存在泄露风险生产环境建议通过受保护的环境变量或密钥管理机制注入。4.2 连接与查询类参数3 个参数类型必填说明restApiEndpointString✅ServiceNow 实例的 REST API 端点例如https://instance.service-now.comqueryModeString✅查询模式二选一Reporting—— 选择应用后拉取该应用下所有表的数据Table—— 直接指定表名拉取数据tableNameString✅要读取数据的 ServiceNow 表名当queryModeReporting时该值会被忽略4.3 取值与输出类参数3 个参数类型必填说明valueTypeString✅返回值的类型Actual—— 返回表中实际存储值默认Display—— 返回表的显示值referenceNameString✅参考名来自 BaseCdapOptionsCDAP 插件通用配置用于标识该数据源outputTxtFilePathPrefixString✅输出文件夹路径 文件名前缀写出时会生成{prefix}-###形式的一组 .txt 文件outputTxtFilePathPrefix的###分片命名规则来自TextIO.write()的标准行为多个 Worker 并行写出时会得到{prefix}-00000-of-00001、{prefix}-00000-of-00002等文件。五、源码级深入参数如何流向 CDAP 插件5.1 参数 → 插件配置 Map 的转换PluginConfigOptionsConverter.java 负责把 Pipeline Options 转成 CDAP 插件认识的参数 Mapreturn ImmutableMap.String, Objectbuilder() .put(ServiceNowConstants.PROPERTY_CLIENT_ID, options.getClientId()) .put(ServiceNowConstants.PROPERTY_CLIENT_SECRET, options.getClientSecret()) .put(ServiceNowConstants.PROPERTY_USER, options.getUser()) .put(ServiceNowConstants.PROPERTY_PASSWORD, options.getPassword()) .put(ServiceNowConstants.PROPERTY_API_ENDPOINT, options.getRestApiEndpoint()) .put(ServiceNowConstants.PROPERTY_QUERY_MODE, options.getQueryMode()) .put(ServiceNowConstants.PROPERTY_TABLE_NAME, options.getTableName()) .put(ServiceNowConstants.PROPERTY_VALUE_TYPE, options.getValueType()) .put(Constants.Reference.REFERENCE_NAME, options.getReferenceName()) .build();可见每个命令行参数都被映射为ServiceNowConstants中的插件属性常量如PROPERTY_CLIENT_ID、PROPERTY_QUERY_MODEreferenceName则映射为 CDAP 通用常量Constants.Reference.REFERENCE_NAME。5.2 完整调用链命令行 --clientId... --queryModeTable ... │ PipelineOptionsFactory 解析 ▼ CdapServiceNowOptions │ PluginConfigOptionsConverter.serviceNowOptionsToParamsMap() ▼ MapString, Object paramsMap │ FormatInputTransform.readFromCdapServiceNow(paramsMap) ▼ ConfigWrapperServiceNowSourceConfig → CdapIO.read() │ 指定 ServiceNowSource 插件 ▼ KVNullWritable, StructuredRecord 数据流 │ MapValues → StructuredRecordUtils.structuredRecordToString ▼ JSON 字符串 → TextIO.write() → {prefix}-###.txt这条链路完整展示了 Beam 示例如何寄生在 CDAP 插件体系之上Beam 负责参数解析、Coder 设置与写出CDAP 插件负责与 ServiceNow REST API 交互。六、输出格式与验证输出内容每行一条 JSON 格式的记录对应 ServiceNow 表的一行数据由structuredRecordToString通过 Gson 序列化生成。输出文件以--outputTxtFilePathPrefix指定的前缀生成的多个分片文件形如{prefix}-00000-of-00001。验证方式管道运行结束后检查输出目录下生成的.txt文件确认每行 JSON 记录中的字段与 ServiceNow 表中查询到的数据一致若valueTypeDisplay字段值应为表定义的显示值而非存储值。七、注意事项与限制参数校验严格9 个必填参数clientId、clientSecret、user、password、restApiEndpoint、queryMode、tableName、valueType、outputTxtFilePathPrefix加上referenceName共 10 项全部标注Validation.Required缺失任一参数管道启动即失败。依赖 CDAP 插件生态本示例依赖io.cdap.plugin.servicenow相关 JAR构建时需确保该依赖在 classpath 中可通过classpath sourceSets.main.runtimeClasspath自动带入。queryMode与tableName联动使用Reporting模式时tableName被忽略选择应用后拉取该应用下所有表的数据使用Table模式时则必须指定具体表名。凭据安全认证信息通过命令行明文传递适合本地演示与测试生产环境应结合安全配置管理。Runner 兼容性示例本身为标准 Beam 批处理管道TextIOMapValues理论上可运行在任意支持批处理的 Beam Runner 上本文以本地 DirectRunner 为例说明。八、进一步探索阅读其他 CDAP 插件示例可对比 Salesforce、Hubspot、Zendesk 等目录下的同构实现加深对CdapIO通用接入模式的理解。深入CdapIO底层其桥接实现位于org.apache.beam.sdk.io.cdap包Java SDK 的 core 相关源码可进一步研究CdapIO.Read如何基于 CDAP 插件构建 Beam 数据源。如果需求是从 ServiceNow 读取数据后再做转换、聚合或写入其他存储如 BigQuery、GCS只需在上述三步流水线中插入对应 Beam Transform 即可本示例是最简化的可扩展骨架。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam Java Kata 实战用 TextIO.read() 从文本文件读取 PCollectionApache Beam Java Kata 实战用 TextIO.read 从文本文件读取 PCollection 本篇技术指南以 Apache Beam 官大数据批处理流处理数据工程Apache Beam CdapIO 实战构建从 ServiceNow 批量拉取数据并落盘 TXT 的管道示例Apache Beam CdapIO 实战构建从 ServiceNow 批量拉取数据并落盘 TXT 的管道示例 本文基于 Apache Beam 仓库中 ex大数据批处理流处理数据工程Apache Beam Kotlin Kata 实战使用 TextIO 从文本文件读取 PCollectionApache Beam Kotlin Kata 实战使用 TextIO 从文本文件读取 PCollection 创建 Beam 管道时最常见的第一步就是从文大数据批处理流处理数据工程上一篇5分钟掌握CompressO免费开源视频图片压缩终极指南下一篇ALVR无线VR串流终极指南彻底告别线缆束缚的完整解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考