Écrire des données depuis Dataflow vers des bases de données
Restez organisé à l'aide des collections
Enregistrez et classez les contenus selon vos préférences.
Pour écrire depuis Dataflow vers des bases de données relationnelles, utilisez le
connecteur d'E/S géré.
L'exemple suivant écrit quelques enregistrements dans une base de données PostgreSQL.
Bien que cet exemple utilise PostgreSQL, la configuration d'autres bases de données compatibles est similaire.
import staticorg.apache.beam.sdk.schemas.Schema.toSchema;importcom.google.common.collect.ImmutableMap;importjava.util.Arrays;importjava.util.List;importjava.util.stream.Stream;importorg.apache.beam.sdk.Pipeline;importorg.apache.beam.sdk.PipelineResult;importorg.apache.beam.sdk.managed.Managed;importorg.apache.beam.sdk.options.Description;importorg.apache.beam.sdk.options.PipelineOptions;importorg.apache.beam.sdk.options.PipelineOptionsFactory;importorg.apache.beam.sdk.schemas.Schema;importorg.apache.beam.sdk.transforms.Create;importorg.apache.beam.sdk.values.Row;publicclassPostgresWrite{privatestaticSchemaINPUT_SCHEMA=Stream.of(Schema.Field.of("id",Schema.FieldType.INT32),Schema.Field.of("name",Schema.FieldType.STRING)).collect(toSchema());privatestaticList<Row>ROWS=Arrays.asList(Row.withSchema(INPUT_SCHEMA).withFieldValue("id",1).withFieldValue("name","John Doe").build(),Row.withSchema(INPUT_SCHEMA).withFieldValue("id",2).withFieldValue("name","Jane Smith").build());publicinterfaceOptionsextendsPipelineOptions{@Description("The JDBC URL of the PostgreSQL database to write to.")StringgetJdbcUrl();voidsetJdbcUrl(Stringvalue);@Description("The PostgresSQL table to write to.")StringgetTable();voidsetTable(Stringvalue);@Description("The username for the PostgreSQL database.")StringgetUsername();voidsetUsername(Stringvalue);@Description("The password for the PostgreSQL database.")StringgetPassword();voidsetPassword(Stringvalue);}publicstaticPipelineResult.Statemain(String[]args){// Parse the pipeline options passed into the application. Example:// --runner=DirectRunner --jdbcUrl=$JDBC_URL --table=$TABLE// --username=$USERNAME --password=$PASSWORD// For more information, see// https://beam.apache.org/documentation/programming-guide/#configuring-pipeline-optionsvaroptions=PipelineOptionsFactory.fromArgs(args).withValidation().as(Options.class);Pipelinepipeline=createPipeline(options);returnpipeline.run().waitUntilFinish();}publicstaticPipelinecreatePipeline(Optionsoptions){// Create configuration parameters for the Managed I/O transform.ImmutableMap<String,Object>config=ImmutableMap.<String,Object>builder().put("jdbc_url",options.getJdbcUrl()).put("location",options.getTable()).put("username",options.getUsername()).put("password",options.getPassword()).build();// Build the pipeline.varpipeline=Pipeline.create(options);pipeline// Create data to write to Postgres..apply(Create.of(ROWS)).setRowSchema(INPUT_SCHEMA)// Write data to a Postgres database using Managed I/O..apply(Managed.write(Managed.POSTGRES).withConfig(config)).getSinglePCollection();returnpipeline;}}
Sauf indication contraire, le contenu de cette page est régi par une licence Creative Commons Attribution 4.0, et les échantillons de code sont régis par une licence Apache 2.0. Pour en savoir plus, consultez les Règles du site Google Developers. Java est une marque déposée d'Oracle et/ou de ses sociétés affiliées.
Dernière mise à jour le 2026/09/10 (UTC).
[[["Facile à comprendre","easyToUnderstand","thumb-up"],["J'ai pu résoudre mon problème","solvedMyProblem","thumb-up"],["Autre","otherUp","thumb-up"]],[["Difficile à comprendre","hardToUnderstand","thumb-down"],["Informations ou exemple de code incorrects","incorrectInformationOrSampleCode","thumb-down"],["Il n'y a pas l'information/les exemples dont j'ai besoin","missingTheInformationSamplesINeed","thumb-down"],["Problème de traduction","translationIssue","thumb-down"],["Autre","otherDown","thumb-down"]],["Dernière mise à jour le 2026/09/10 (UTC)."],[],[]]