Geotools et performance

Il y a quelques semaines, j'ai présenté un exemple d'utilisation de la librairie Geotools pour la manipulation de données géographiques.

Après avoir éprouvé ce code face à des volumétries importantes, j'ai pu constater quelques problèmes de montées en charge de la JVM, tels qu'illustrés ici :

Face à un jeu de test d'environ 200Mo de fichiers Shapefile (représentant approximativement 700.000 lignes en base), on voit clairement les pics de charge dus à l'ajout en mémoire du contenu total des fichiers avant injection en base de données.

Je vais donc ici présenter une façon d'améliorer le code décrit précédemment afin de mieux gérer la mémoire utilisée.

Geotools fournit deux objets très utiles pour gérer les données utilisées sous la forme de flux :

  • Dans le cas de l'injection en base de données provenant de Shapefiles : FeatureReader<SimpleFeatureType, SimpleFeature> qui permet de lire en continu les features présentes dans un fichier
  • Dans le cas de l'extraction des données d'une base vers un fichier Shapefile : Query.setMaxFeatures(int) et Query.setStartIndex(int) permettent la pagination des requêtes SQL.

Voici donc comment les utiliser :

private static final int MAX_MEMORY_FEATURES = 15000;

private static final String POSTGIS_TABLENAME = "MY_TABLE";

private static GeoProperties props = GeoProperties.getInstance();

private static CamelProperties camelprops = CamelProperties.getInstance();

private static ShapefileDataStoreFactory shpFactory = new ShapefileDataStoreFactory();

private static FeatureTypeFactoryImpl factory = new FeatureTypeFactoryImpl();

private static DataStore pgStore;

private SimpleFeatureType schema;

private static Log LOGGER = LogFactory.getLog("camel");

public GeoToPostGISClient() throws Exception {
 try {
  Logging.GEOTOOLS.setLoggerFactory("org.geotools.util.logging.Log4JLoggerFactory");
 } catch (Exception e) {
  LOGGER.warn("Log factory not found for GeoTools : 'org.geotools.util.logging.Log4JLoggerFactory'");
 }
 // Ouvrir une connexion vers la base PostGIS
 getDataStore();
}

private void getDataStore() throws Exception {
 if (pgStore == null) {
  PostgisDataStoreFactory pgFactory = new PostgisDataStoreFactory();
  Map<String, String> jdbcparams = new HashMap<String, String>();
  jdbcparams.put(PostgisDataStoreFactory.DBTYPE.key, "postgis");
  jdbcparams.put(PostgisDataStoreFactory.HOST.key, props.getProperty(GeoProperties.DB_HOST));
  jdbcparams.put(PostgisDataStoreFactory.PORT.key, props.getProperty(GeoProperties.DB_PORT));
  jdbcparams.put(PostgisDataStoreFactory.SCHEMA.key, props.getProperty(GeoProperties.DB_SCHEMA));
  jdbcparams.put(PostgisDataStoreFactory.DATABASE.key, props.getProperty(GeoProperties.DB_NAME));
  jdbcparams.put(PostgisDataStoreFactory.USER.key, props.getProperty(GeoProperties.DB_USER));
  jdbcparams.put(PostgisDataStoreFactory.PASSWD.key, props.getProperty(GeoProperties.DB_PWD));
  pgStore = pgFactory.createDataStore(jdbcparams);
 }
}

/**
 * Insert all specified shapefiles in Postgre
 * 
 * @param shapefilePaths files
 * @throws IOException all
 */
public void insertShpIntoDb(List<String> shapefilePaths) throws Exception {
 Iterator<String> iterator = shapefilePaths.iterator();
 String path = null;
 while (iterator.hasNext()) {
  path = iterator.next();
  LOGGER.info("Inserting : " + path);

  Map<String, Object> shpparams = new HashMap<String, Object>();
  shpparams.put("url", "file://" + path);
  FileDataStore shpStore = (FileDataStore) shpFactory.createDataStore(shpparams);

  if (schema == null) {
   LOGGER.info("Create schema");
   // Copy schema and change name in order to refer to the same
   // global schema for all files
   SimpleFeatureType originalSchema = shpStore.getSchema();
   Name originalName = originalSchema.getName();
   NameImpl theName = new NameImpl(originalName.getNamespaceURI(), originalName.getSeparator(), POSTGIS_TABLENAME);
   schema = factory.createSimpleFeatureType(theName, originalSchema.getAttributeDescriptors(), originalSchema.getGeometryDescriptor(),
     originalSchema.isAbstract(), originalSchema.getRestrictions(), originalSchema.getSuper(), originalSchema.getDescription());
   pgStore.createSchema(schema);
  }

  // Ajout des objets du shapefile dans la table PostGIS
  // Query.FIDS : To request only the feature IDs with no content
  int totalSHPentries = shpStore.getFeatureSource().getCount(Query.FIDS);
  int insertedEntries = 0;
  SimpleFeatureStore featureStore = (SimpleFeatureStore) pgStore.getFeatureSource(POSTGIS_TABLENAME);
  DefaultTransaction transaction = null;
  FeatureReader<SimpleFeatureType, SimpleFeature> featureReader = shpStore.getFeatureReader();
  SimpleFeatureCollection features = new DefaultFeatureCollection(null, featureReader.getFeatureType());
  while (featureReader.hasNext()) {
   features.add(featureReader.next());
   if (features.size() == MAX_MEMORY_FEATURES || !featureReader.hasNext()) {
    transaction = new DefaultTransaction("bulk");
    featureStore.setTransaction(transaction);
    try {
     LOGGER.info("Inserting features " + insertedEntries + " to " + (insertedEntries + features.size()) + " from "
       + totalSHPentries);
     featureStore.addFeatures(features);
     transaction.commit();
     insertedEntries += features.size();
     features = new DefaultFeatureCollection(null, featureReader.getFeatureType());
     // To avoid memory leaks
     System.gc();
    } catch (Exception problem) {
     LOGGER.error(problem.getMessage(), problem);
     transaction.rollback();
     break;
    } finally {
     transaction.close();
    }
   }
  }
  featureReader.close();
  features = null;
  // To avoid memory leaks
  System.gc();

  shpStore.dispose();
  LOGGER.info("End insert");
 }
 extractFromDb();
}

/**
 * Extracts local data from postgis DB
 * 
 * @throws IOException all
 */
public void extractFromDb() throws IOException {
 // Faire une requête spatiale dans la base
 SimpleFeatureCollection filteredFeatures = null;

 String destFolder = camelprops.getProperty(CamelProperties.CAMEL_WORK_DIR) + "/shp/";

 LOGGER.info("Extracting data");
 for (Object dep : ReferentielDepartement.getDepartements()) {
  try {
   // Check data presence in DB
   Filter deptFilter = CQL.toFilter("DPT_NUM = '" + dep + "'");
   Query qCount = new Query(pgStore.getTypeNames()[0], deptFilter);
   int count = pgStore.getFeatureSource(POSTGIS_TABLENAME).getCount(qCount);

   if (count > 0) {
    // Écrire le résultat dans un fichier shapefile
    Map<String, String> destshpparams = new HashMap<String, String>();
    String destinationSchemaName = ShapeFileNamesRules.computeNameFor((String) dep);
    destshpparams.put("url", "file://" + destFolder + destinationSchemaName + ".shp");
    DataStore destShpStore = shpFactory.createNewDataStore(destshpparams);

    // duplicate existing schema to create destination's one
    Name originalName = schema.getName();
    NameImpl theName = new NameImpl(originalName.getNamespaceURI(), originalName.getSeparator(), destinationSchemaName);
    SimpleFeatureType destschema = factory
      .createSimpleFeatureType(theName, schema.getAttributeDescriptors(), schema.getGeometryDescriptor(), schema.isAbstract(),
        schema.getRestrictions(), schema.getSuper(), schema.getDescription());
    destShpStore.createSchema(destschema);

    // destination store
    SimpleFeatureStore destFeatureStore = (SimpleFeatureStore) destShpStore.getFeatureSource(destinationSchemaName);

    // Query DB
    int extractedData = 0;
    Query q = null;
    while (extractedData < count) {
     q = new DefaultQuery(pgStore.getTypeNames()[0], deptFilter, MAX_MEMORY_FEATURES, null, "extractQuery");
     q.setStartIndex(extractedData);
     filteredFeatures = pgStore.getFeatureSource(POSTGIS_TABLENAME).getFeatures(q);
     if (filteredFeatures != null && filteredFeatures.size() > 0) {
      LOGGER.info("Extracting features " + extractedData + " to " + filteredFeatures.size() + " from " + count + " for " + dep);
       destFeatureStore.addFeatures(filteredFeatures);
     }
     extractedData += MAX_MEMORY_FEATURES;
    }

    // Fermer les connections et les fichiers
    LOGGER.info("End extract");
    destShpStore.dispose();
   }
  } catch (CQLException e) {
   LOGGER.error(e.getMessage(), e);
  }
 }

 // Write done file for Camel
 File done = new File(destFolder + "done");
 done.createNewFile();
 done = null;
}

Avec ce code, les données sont insérées en base ou extraites par lot de 15000, ce qui permet de limiter la surcharge de la JVM. Voici le résultat :

On voit distinctement les montées en charge pour chaque fichier traité, mais le tout est lissé et limité autour de 30Mo, ce qui permet de plus facilement calibrer l'environnement final de l'application. Un exemple de log produit :

780617 INFO  camel  - Inserting : C:\test\shp\RPG_2010_064.shp
781164 INFO  camel  - Inserting features 0 to 15000 from 120834
813128 INFO  camel  - Inserting features 15000 to 30000 from 120834
845466 INFO  camel  - Inserting features 30000 to 45000 from 120834
875914 INFO  camel  - Inserting features 45000 to 60000 from 120834
907550 INFO  camel  - Inserting features 60000 to 75000 from 120834
937530 INFO  camel  - Inserting features 75000 to 90000 from 120834
968837 INFO  camel  - Inserting features 90000 to 105000 from 120834
998942 INFO  camel  - Inserting features 105000 to 120000 from 120834
1028405 INFO  camel  - Inserting features 120000 to 120834 from 120834
1030327 INFO  camel  - End insert
1030327 INFO  camel  - Inserting : C:\test\shp\RPG_2010_080.shp
1030733 INFO  camel  - Inserting features 0 to 15000 from 86010
1056963 INFO  camel  - Inserting features 15000 to 30000 from 86010
1084584 INFO  camel  - Inserting features 30000 to 45000 from 86010
1111438 INFO  camel  - Inserting features 45000 to 60000 from 86010
1136965 INFO  camel  - Inserting features 60000 to 75000 from 86010
1161602 INFO  camel  - Inserting features 75000 to 86010 from 86010
1180505 INFO  camel  - End insert

Sources :


Fichier(s) joint(s) :



Apache Camel et performance

Après avoir bien pris en main Camel, je dois dire que je suis tout à fait satisfait des possibilités offertes (voir mon précédent article). Cependant, j'ai remarqué quelques comportements peu adaptés au traitement des volumétries importantes de données.

En effet, il semblerait que certains composants et/ou outils proposés par défaut causent des pics de mémoire voir même systématiquement des crash OutOfMemory. Je vais donc essayer ici de les décrire afin de permettre aux futurs développeurs de les prendre en compte dès le départ. Je précise tout de même que j'utilise actuellement la version 2.8.3.

Marshalers et OutOfMemory

Voici un exemple de code simpliste situant le problème :

from("file://...").unmarshal().zip().to("file://...");

Ici, le but est de dezipper un fichier pour écrire son contenu vers la destination spécifiée. Utilisez un fichier de 500Mo par exemple et l'application devrait planter assez rapidement!

En effet, la méthode unmarshall().zip() va charger en mémoire (!) un DeflaterOutputStream correspondant au contenu total du fichier AVANT de l'écrire. Qui plus est, un fichier plus modeste (quelques Mo), ne causant pas de crash, restera chargé en mémoire jusqu'à l'arrivé du garbage collector.

L'intérêt de conserver en mémoire ces informations peut être compréhensible dans le cas du traitement de petits fichiers, mais dès lors que l'on pousse l'implémentation vers des volumétries importantes, il vaut mieux privilégier une autre stratégie (Camel ne prévoyant à priori pas la possibilité d'écrire le flux de données au fur et à mesure de se décompression). Un code plus robuste ressemblerait à :

from("file://...").convertBodyTo(File.class).to("bean://com.MyDezipperBean");

Qui se contente de rediriger le fichier vers un bean personnalisé qui permettra de mieux gérer les flux de données, notamment grâce au classique java.util.zip.*

Notez qu'il en va de même pour le CSV. L'appel à la méthode unmarshal().csv() produira le même comportement. Préférez donc l'utilisation de org.apache.commons.csv.CSVParser par l'intermédiaire d'un bean, beaucoup plus à-même de gérer les fichiers CSV volumineux.

Split et montées en charge

Imaginons la nécessité de découper un gros fichier CSV entrant en plusieurs plus petits (quelques soient les conditions et méthodes employées). Une implémentation évidente ressemblerait à :

from("file://...").convertBodyTo(File.class)
  .split().method("com.services.CsvService","splitDatas")
  .to("file://...");

com.services.CsvService est un bean recevant le fichier CSV initial, le parcourant pour découper le contenu selon X règles métier puis retournant une liste de contenus CSV destinés à être écrits dans des fichiers séparés.

En observant de près le comportement de l'application, on aura vite fait de se rendre compte de la rapide montée en charge due essentiellement au stockage en mémoire des informations CSV (lors de leur transfert sous forme de message du bean vers l'endpoint fichier). Une fois de plus donc, le cas du traitement de fichiers nombreux et/ou conséquents peut vite devenir problématique. Une solution de contournement peut consister à faire écrire directement au bean les fichiers (pour ne pas avoir à les re-transporter sur Camel) puis éventuellement de faire transiter simplement leur chemin si besoin.

Voici donc les problématiques que j'ai rencontrées à ce jour. Après avoir été plusieurs fois agréablement surpris par l'intelligence des composants et des diverses implémentations de EIP proposées par Camel, je dois avouer que j'ai quelque peu été dérouté par ces comportements.

Hope this helps!


Fichier(s) joint(s) :



Séparer les logs des modules avec Log4j

Dans des applications complexes, il peut être intéressant de séparer les logs des différents modules dans des fichiers distincts afin d'améliorer leur lecture.

Je vais présenter ici un exemple de fichier log4j.properties utilisé pour réguler la destination et le niveau de sortie des logs selon leur type, dans une application basée sur Camel et la librairie GeoTools :

# Set root logger level to DEBUG and its direct appenders to camel and console.
log4j.rootLogger=DEBUG, camel, console

log4j.logger.org.geotools=DEBUG, geoToolsAppender
# Set GeoTools JDBC appender (for SQL showing)
log4j.logger.org.geotools.jdbc=DEBUG, geoToolsJDBCAppender
# To avoid duplication in camel's logs
log4j.additivity.org.geotools.jdbc=false
log4j.additivity.org.geotools=false

log4j.appender.console=org.apache.log4j.ConsoleAppender
log4j.appender.console.layout=org.apache.log4j.PatternLayout
log4j.appender.console.layout.ConversionPattern=%-4r [%t] %-5p %c %x - %m%n
log4j.appender.console.Threshold=INFO

log4j.appender.camel=org.apache.log4j.FileAppender
log4j.appender.camel.File=camel.log
log4j.appender.camel.layout=org.apache.log4j.PatternLayout
log4j.appender.camel.layout.ConversionPattern=%-4r [%t] %-5p %c %x - %m%n
log4j.appender.camel.Threshold=DEBUG

log4j.appender.geoToolsAppender=org.apache.log4j.FileAppender
log4j.appender.geoToolsAppender.File=geotools.log
log4j.appender.geoToolsAppender.layout=org.apache.log4j.PatternLayout
log4j.appender.geoToolsAppender.layout.ConversionPattern=%-4r [%t] %-5p %c %x - %m%n
log4j.appender.geoToolsAppender.Threshold=DEBUG

log4j.appender.geoToolsJDBCAppender=org.apache.log4j.FileAppender
log4j.appender.geoToolsJDBCAppender.File=geotoolsjdbc.log
log4j.appender.geoToolsJDBCAppender.layout=org.apache.log4j.PatternLayout
log4j.appender.geoToolsJDBCAppender.layout.ConversionPattern=%-4r [%t] %-5p %c %x - %m%n
log4j.appender.geoToolsJDBCAppender.Threshold=DEBUG

# To avoid annoying polling logs
log4j.logger.org.apache.camel.component.file.FileConsumer=ERROR

Dans cet exemple, quatre logs sont mis en place :

  • console : pour diriger certaines informations vers la console Eclipse.
  • camel : pour créer un log spécifique utilisé par l'application Camel (appelé par org.apache.commons.logging.Log LOGGER = org.apache.commons.logging.LogFactory.getLog("camel");).
  • geoToolsAppender : un logger spécifique utilisé par un module de l'application (GeoTools) : toutes les sorties issues des classes à l'intérieur des packages org.geotools.* utiliseront l'appender geoToolsAppender
  • geoToolsJDBCAppender : même chose que le point précédent, spécifique au package org.geotools.jdbc (pour faire apparaître notamment des requêtes SQL)

Les instructions log4j.additivity... servent à indiquer à Log4j de ne pas dupliquer les sorties des logs en question dans le log principal (rootLogger dirigé vers "camel").

Voilà tout!

Sources :


Fichier(s) joint(s) :



Manipuler des données géographiques avec GeoTools

Dans le monde des Systèmes d'Informations Géographique (SIG), il est assez courant de devoir manipuler des données au format Shapefile (ensemble de fichiers standardisés pour la représentation de cartographie) afin de présenter différents types d'informations.

Pour ce faire, le plus simple est encore d'utiliser comme intermédiaire une base Postgre dédiée, autrement nommée PostGIS, spécifique au traitement de ce type de données.

La solution la plus courante est d'utiliser les outils en ligne de commande shp2pgsql ou ogr2ogr qui permettent de créer un fichier SQL à partir des Shapefile et éventuellement de le jouer directement en base. Cependant, le but de ce billet est de présenter comment il est possible en Java d'extraire des informations et/ou ré-organiser un lot de fichiers Shapefile. L'utilisation d'une base de données intermédiaire a pour vocation de résoudre des problèmes de performance liés à la manipulation directe des fichiers.

Voici un exemple de code à utiliser :

public class GeoToPostGISClient {

 private static final String POSTGIS_TABLENAME = "MY_TABLE";

 private static GeoProperties props = GeoProperties.getInstance();

 private static ShapefileDataStoreFactory shpFactory = new ShapefileDataStoreFactory();

 private static FeatureTypeFactoryImpl factory = new FeatureTypeFactoryImpl();

 private static JDBCDataStore pgStore;

 private SimpleFeatureType schema;

 public GeoToPostGISClient() throws IOException {
  // Ouvrir une connexion vers la base PostGIS
  if (pgStore == null) {
   PostgisNGDataStoreFactory pgFactory = new PostgisNGDataStoreFactory();
   Map<String, String> jdbcparams = new HashMap<String, String>();
   jdbcparams.put(PostgisNGDataStoreFactory.DBTYPE.key, "postgis");
   jdbcparams.put(PostgisNGDataStoreFactory.HOST.key, props.getProperty(GeoProperties.DB_HOST));
   jdbcparams.put(PostgisNGDataStoreFactory.PORT.key, props.getProperty(GeoProperties.DB_PORT));
   jdbcparams.put(PostgisNGDataStoreFactory.SCHEMA.key, props.getProperty(GeoProperties.DB_SCHEMA));
   jdbcparams.put(PostgisNGDataStoreFactory.DATABASE.key, props.getProperty(GeoProperties.DB_NAME));
   jdbcparams.put(PostgisNGDataStoreFactory.USER.key, props.getProperty(GeoProperties.DB_USER));
   jdbcparams.put(PostgisNGDataStoreFactory.PASSWD.key, props.getProperty(GeoProperties.DB_PWD));
   pgStore = pgFactory.createDataStore(jdbcparams);
  }
 }

 /**
  * Insert all specified shapefiles in Postgre
  * 
  * @param shapefilePaths files
  * @throws IOException all
  */
 public void insertShpIntoDb(List<String> shapefilePaths) throws IOException {
  Iterator<String> iterator = shapefilePaths.iterator();
  String path = null;
  while (iterator.hasNext()) {
   path = iterator.next();

   Map<String, Object> shpparams = new HashMap<String, Object>();
   shpparams.put("url", "file://" + path);
   // create indexes only for last file (performance issue)
   FileDataStore shpStore = (FileDataStore) shpFactory.createDataStore(shpparams);
   SimpleFeatureCollection features = shpStore.getFeatureSource().getFeatures();

   if (schema == null) {
    // Copy schema and change name in order to refer to the same
    // global schema for all files
    SimpleFeatureType originalSchema = shpStore.getSchema();
    Name originalName = originalSchema.getName();
    NameImpl theName = new NameImpl(originalName.getNamespaceURI(), originalName.getSeparator(), POSTGIS_TABLENAME);
    schema = factory.createSimpleFeatureType(theName, originalSchema.getAttributeDescriptors(), originalSchema.getGeometryDescriptor(),
      originalSchema.isAbstract(), originalSchema.getRestrictions(), originalSchema.getSuper(), originalSchema.getDescription());
    pgStore.createSchema(schema);
   }
   SimpleFeatureStore featureStore = (SimpleFeatureStore) pgStore.getFeatureSource(POSTGIS_TABLENAME);

   // Ajout des objets du shapefile dans la table PostGIS
   DefaultTransaction transaction = new DefaultTransaction("create");
   featureStore.setTransaction(transaction);
   try {
    featureStore.addFeatures(features);
    transaction.commit();
   } catch (Exception problem) {
    LOGGER.error(problem.getMessage(), problem);
    transaction.rollback();
   } finally {
    transaction.close();
   }
   shpStore.dispose();
  }
  extractFromDb();
 }

 /**
  * Extracts local data from postgis DB
  * 
  * @throws IOException all
  */
 public void extractFromDb() throws IOException {
  // Faire une requête spatiale dans la base
  ContentFeatureCollection filteredFeatures = null;

  String destFolder = "/shp/";

  for (Object dep : ReferentielDepartement.getDepartements()) {
   try {
    filteredFeatures = pgStore.getFeatureSource(POSTGIS_TABLENAME).getFeatures(CQL.toFilter("DPT_NUM = '" + dep + "'"));
   } catch (CQLException e) {
    LOGGER.error(e.getMessage(), e);
   }
   if (filteredFeatures != null && filteredFeatures.size() > 0) {
    // Écrire le résultat dans un fichier shapefile
    Map<String, String> destshpparams = new HashMap<String, String>();
    SimpleDateFormat formatter = new SimpleDateFormat("yyyyMMdd");
    String destinationSchemaName = "MySchema_" + dep;
    destshpparams.put("url", "file://" + destFolder + destinationSchemaName + "_" + formatter.format(new Date()) + ".shp");
    DataStore destShpStore = shpFactory.createNewDataStore(destshpparams);

    // duplicate existing schema to create destination's one
    Name originalName = schema.getName();
    NameImpl theName = new NameImpl(originalName.getNamespaceURI(), originalName.getSeparator(), destinationSchemaName);
    SimpleFeatureType destschema = factory.createSimpleFeatureType(theName, schema.getAttributeDescriptors(),
      schema.getGeometryDescriptor(), schema.isAbstract(), schema.getRestrictions(), schema.getSuper(), schema.getDescription());
    destShpStore.createSchema(destschema);

    SimpleFeatureStore destFeatureStore = (SimpleFeatureStore) destShpStore.getFeatureSource(destinationSchemaName);
    destFeatureStore.addFeatures(filteredFeatures);

    // Fermer les connections et les fichiers
    destShpStore.dispose();
   }
  }
 }
}

Avec ce type de code, il est possible d'extraire une nouvelle cartographie spécifique (découpage selon la variable DPT_NUM) à partir d'un lot de données source.

Pour une mise en place plus rapide, voici les dépendances nécessaires (pom.xml) :

...
<repositories>
 <repository>
  <id>osgeo</id>
  <name>Open Source Geospatial Foundation Repository</name>
  <url>http://download.osgeo.org/webdav/geotools/</url>
 </repository>
</repositories>
...
<dependencies>
 <!-- Geo Tools -->
 <dependency>
  <groupId>org.geotools</groupId>
  <artifactId>gt-shapefile</artifactId>
  <version>8.0-M4</version>
 </dependency>
 <dependency>
  <groupId>org.geotools.jdbc</groupId>
  <artifactId>gt-jdbc-postgis</artifactId>
  <version>8.0-M4</version>
 </dependency>
 <dependency>
  <groupId>org.geotools</groupId>
  <artifactId>gt-cql</artifactId>
  <version>8.0-M4</version>
 </dependency>
</dependencies>

Voilà tout, bon courage!
HTH


Fichier(s) joint(s) :



Apache Camel par l'exemple

Avec cet article j'ai décidé d'entrer directement dans le vif du sujet...

S'il fallait présenter rapidement Camel, on pourrait dire qu'il s'agit d'une plateforme d'intégration d'application, basée sur un système d'échange de messages et dont le but est de fournir une implémentation des grands patrons d'intégrations en entreprise (facilitant la communication inter-applications). Pour ne pas plagier ou paraphraser, voici deux articles intéressants présentant ces patrons : le premier sur le site de Novedia et le second, plus détaillé chez Soat. Pour continuer sur une présentation plus spécifique de Camel, voici un premier billet écrit sur le blog d'Octo en enfin une présentation complète par un des co-auteurs du livre "Camel In Action", Jonathan Anstey.

Mon but ici est donc de fournir un exemple de mise en place d'un "bus" Camel pour créer un flux applicatif.

Le scénario est le suivant : on doit récupérer une archive zippée sur un serveur FTP distant, la décompresser et traiter son contenu en fonction de son type : les fichiers CSV doivent être segmentés selon une règle métier puis re-zippés unitairement et les autres types de fichiers sont envoyés à un script shell. Entre-temps, les données sont triées et validées. Celles qui sont invalides sont déposées séparément dans un répertoire spécifique.

De manière plus illustrée :

En jaune sont représentés les composants intégrés à Camel : FTP, ZIP, EXEC et CSV
En rouge sont illustrés les patrons d'intégration implémentés : Split, Enricher, Router, Filter, Sort, Recipient list, Validate.
En blanc sont indiqués les beans/services personnalisés ajoutés.

Et maintenant le plus intéressant, le code pour la mise en place des routes (les commentaires décrivent tout son fonctionnement) :

import java.util.Comparator;
import java.util.List;

import org.apache.camel.Exchange;
import org.apache.camel.Processor;
import org.apache.camel.builder.RouteBuilder;
import org.apache.camel.dataformat.csv.CsvDataFormat;
import org.apache.camel.language.bean.BeanLanguage;
import org.apache.camel.processor.validation.PredicateValidationException;

import CamelProperties;
import ReferentielDepartement;

/**
 * Main class to creates Camel routes
 * 
 * @author pe.faidherbe
 *
 */
public class MyRoutesBuilder extends RouteBuilder {
 
 // Stores archive's file name retrieved from remote server
 private static String downloadedFileName;
 
 // Indicates Csv column's name used to sort datas
 private static String csvIdentityColumn;
 
 // Indicates Csv column's number used to sort datas
 private static Integer csvIdentityColumnNumber;
 
 // Indicates length of sorting data used as data identifier
 private static int csvIdentityLength;
 
 // General properties
 private static final CamelProperties camelProps = CamelProperties.getInstance();
 
 // Csv data delimiter
 private static final String csvDelimiter = camelProps.getProperty(CamelProperties.CAMEL_DATA_SEP);
 
 // Prefix of data file (used for routing)
 private static final String NAT_FILE_PREFIX = camelProps.getProperty(CamelProperties.CAMEL_DATA_FILE_NAT_PREFIX);
 
 // Prefix of data file (used for routing)
 private static final String S2_FILE_PREFIX = camelProps.getProperty(CamelProperties.CAMEL_DATA_FILE_S2_PREFIX);
 
 // Prefix of data file (used for routing)
 private static final String ILOT_FILE_PREFIX = camelProps.getProperty(CamelProperties.CAMEL_DATA_FILE_ILOT_PREFIX);
 
 /**
  * Getter used by Camel
  * @return remote File Name
  */
 public String getDownloadedFileName() {
  return downloadedFileName;
 }
 
 /**
  * Getter used by Camel
  * @return csv identity column
  */
 public String getCsvIdentColumn() {
  return csvIdentityColumn;
 }
 
 /**
  * Getter used by Camel
  * @return csv identity column num
  */
 public Integer getCsvIdentityColumnNumber() {
  return csvIdentityColumnNumber;
 }
 
 /**
  * Getter used by Camel
  * @return csv identity data length
  */
 public int getCsvIdentLength() {
  return csvIdentityLength;
 }
 
 /**
  * Getter used by Camel
  * @return csv delimiter to use
  */
 public String getCsvDelimiter() {
  return csvDelimiter;
 }
 
 /**
  * Used by first Camel route to "persist" informations about remote file
  * later passed as parameters for second route
  */
 private void setMyContext(String downloadedArchive) {
  downloadedFileName = downloadedArchive.substring(0, downloadedArchive.indexOf("."));
  if(downloadedFileName.startsWith(S2_FILE_PREFIX)) {
   csvIdentityColumn = camelProps.getProperty(CamelProperties.CAMEL_DATA_IDENT_COL_S2);
   csvIdentityLength = Integer.parseInt(camelProps.getProperty(CamelProperties.CAMEL_DATA_IDENT_S2_LENGTH));
   csvIdentityColumnNumber = Integer.parseInt(camelProps.getProperty(CamelProperties.CAMEL_DATA_IDENT_COL_S2_NUM));
  } else if(downloadedFileName.startsWith(NAT_FILE_PREFIX)) {
   csvIdentityColumn = camelProps.getProperty(CamelProperties.CAMEL_DATA_IDENT_COL_NAT);
   csvIdentityLength = Integer.parseInt(camelProps.getProperty(CamelProperties.CAMEL_DATA_IDENT_NAT_LENGTH));
   csvIdentityColumnNumber = Integer.parseInt(camelProps.getProperty(CamelProperties.CAMEL_DATA_IDENT_COL_NAT_NUM));
  }
 }

 /**
  * @see org.apache.camel.builder.RouteBuilder#configure()
  */
 @Override
 public void configure() throws Exception {
  // Props
  String camelWorkDir = camelProps.getProperty(CamelProperties.CAMEL_WORK_DIR);
  
  // Csv comparator
  CsvSorter sorter = new CsvSorter();
  
  // dead Letter Channel
  errorHandler(deadLetterChannel("log:camel"));
  
  // not validated messages go to particular error folder
  onException(PredicateValidationException.class).handled(true)
   .to("file://C:/test2/csverror?fileName=${header:downloadedFileName}_${header:territoire}_${date:now:yyyyMMddHHmmss}.csv")
   .log("Validation Error : ${exception.message}").end(); 
  
  /*
   * Download
   */
  from("ftp://"+camelProps.getProperty(CamelProperties.FTP_USER_PROP)
    + "@"
    + camelProps.getProperty(CamelProperties.FTP_HOST)
    + "?password="
    + camelProps.getProperty(CamelProperties.FTP_USER_PWD)
    + "&binary=true&noop=true&disconnect=true"
    // Poll every X sec
    + "&consumer.delay=" + camelProps.getProperty(CamelProperties.FTP_POLL_TIME_MS)
    // Specify temp destination for performance issue (not loaded in memory)
    + "&localWorkDirectory="+camelWorkDir)
   .log("Unzipping : ${file:name}")
   // Read as zip file
   .marshal().zip()
   // Unzip in memory 
   .unmarshal().zip()
   // Send to bean to extract entries
   .split().method("ZipService","unzipFile")
    .log("Extracted : ${header:entryName}")
   // Write each file
   .to("file://"+camelWorkDir+"?fileName=${header:entryName}")
   .process(new Processor() {
    @Override
    public void process(Exchange e) throws Exception {
     // Set informations on how to treat latter data
     setMyContext((String) e.getIn().getHeader("CamelFileName"));
    }
   })
  .end();
  
  /*
   *  ROUTER
   */
  from("file://"+camelWorkDir).id("routerRoute").log("Routing start")
   // No autostart to avoid polling "undesired" file (not previously retrieved from ftp)
   //.noAutoStartup()
   // Enrich file polling with context informations
   .enrich("direct:contextEnricher")
   .choice()
    .when(header("downloadedFileName").startsWith(S2_FILE_PREFIX))
     // CSV data, going to split
     .to("direct:surfaces")
    .when(header("downloadedFileName").startsWith(ILOT_FILE_PREFIX))
     // Geo data, go to DB
     .to("direct:ilots")
    .when(header("downloadedFileName").startsWith(NAT_FILE_PREFIX))
     // CSV data, going to split
     .to("direct:national")
    .otherwise()
     .log("Fichier ${file:name} non pris en charge!")
    .end();
  
  /*
   * Content enricher
   */
  from("direct:contextEnricher")
   .setHeader("downloadedFileName", BeanLanguage.bean(getClass(), "getDownloadedFileName"))
   .setHeader("csvIdentityColumn", BeanLanguage.bean(getClass(), "getCsvIdentColumn"))
   .setHeader("csvIdentityLength", BeanLanguage.bean(getClass(), "getCsvIdentLength"))
   .setHeader("csvDelimiter", BeanLanguage.bean(getClass(), "getCsvDelimiter"));
  
  /*
   * Manage geo data
   */
  from("direct:ilots")
   // Manage only SHP files
   .filter(header("entryName").endsWith("shp"))
   .to("file://C:/test2?fileName=${header:entryName}")
   // RecipientList is needed because route is computed at runtime (because of dynamic parameters)
   .recipientList(simple("exec:C:/test2/cmd/ogr2ogr.bat?args=${header:entryName}&workingDir=C:/test2/cmd/&useStderrOnEmptyStdout=true"))
    // Convert to String because cmd return is InputStream
    .convertBodyTo(String.class)
    .log("Command return : ${body} , error : ${header:exec_stderr}")
   .end();
  
  /*
   * Split CSV
   */

  // Used to customized CSV separator
  CsvDataFormat csvFormat = new CsvDataFormat();
  csvFormat.setDelimiter(csvDelimiter);
  
  from("direct:surfaces").convertBodyTo(String.class).unmarshal(csvFormat)
   .sort(body(), sorter)
   // Split CSV content
   // Bean uses StringBuilder and writes CSV content in order to get better performance than
   // creating a lot of List<Map<String, Object>> handled by camel's csv marshaler
   .split().method("CsvService","splitDatas")
   // Business check : is "territoire" a valid data?
   .validate(header("territoire").in(ReferentielDepartement.getDepartements()))
   // Write each data to a proper file
   .to("file://C:/test2/splitted?fileName=Territorial_${header:territoire}_${date:now:yyyyMMdd}.csv")
   .log("Written CSV file for : ${header:territoire}");
  
  from("direct:national").convertBodyTo(String.class).unmarshal(csvFormat)
   .sort(body(), sorter)
   // Split CSV content
   .split().method("CsvService","splitDatas")
   .validate(header("territoire").in(ReferentielDepartement.getDepartements()))
   // Write each data to a proper file
   .to("file://C:/test2/splitted?fileName=${file:onlyname.noext}_${header:territoire}_${date:now:yyyyMMddHHmmss}.csv")
   .log("Written CSV file for : ${header:territoire}");
  
  
  /*
   * Zip final files
   */
  // To handle transformation from camel's DeflaterOutputStream to traditional ZipOutputStream
  CustomizedZipDataFormat zipFormat = new CustomizedZipDataFormat();
  from("file://C:/test2/splitted").marshal(zipFormat)
   .to("file://C:/test2?fileName=${file:onlyname.noext}.zip")
   .log("Zipped : ${file:onlyname.noext}");
 }
 
 class CsvSorter implements Comparator<List<String>> {
  @Override
  public int compare(List<String> o1, List<String> o2) {
   int result = 0;
   // Do not treat first line (headers)
   if(!o1.contains(csvIdentityColumn) && !o2.contains(csvIdentityColumn)) {
    String ccom1 = o1.get(csvIdentityColumnNumber).substring(0, csvIdentityLength);
    String ccom2 = o2.get(csvIdentityColumnNumber).substring(0, csvIdentityLength);
    result = ccom1.compareTo(ccom2);
   }
   return result;
  }
 }
}

Pour ce qui est du service permettant de segmenter les informations CSV, voici son squelette :

public class CsvService {
 
 /**
  * Receives csv informations and split it out to multiple messages
  * Received datas are pre-ordered
  * @param headers in-message headers
  * @param body in-message body (unmarshalled csv content)
  * @return messages containing csv info to be written
  */
 public List<Message> splitDatas(@Headers Map<String, Object> headers,
   @Body List<List<String>> body) {
  List<Message> answer = new ArrayList<Message>();
  
  // headers
  List<String> csvHeaders = body.get(0);
  
  String csvDelimiter = (String) headers.get("csvDelimiter");
  
  // data
  List<List<String>> datas = body.subList(1, body.size());

  (... sort routine ...)

  return answer;
 }
}

Pour ce qui est du service permettant de dézipper l'archive :

public class ZipService {

 /**
  * Splits in message to multiple messages for entries
  * @param headers in headers
  * @param body in body
  * @return one message per zip entry
  */
 public List<Message> unzipFile(@Headers Map<String, Object> headers,
   @Body Object body) {
  List<Message> answer = new ArrayList<Message>();
  try {
   ZipInputStream zis = new ZipInputStream(new ByteArrayInputStream(
     (byte[]) body));
   ZipEntry ze = null;
   String entryName = "";
   String unzippedFiles = "";
   while ((ze = zis.getNextEntry()) != null) {
    entryName = ze.getName();
    unzippedFiles += entryName + ",";
    ByteArrayOutputStream out = new ByteArrayOutputStream();
    for (int c = zis.read(); c != -1; c = zis.read()) {
     out.write(c);
    }
    zis.closeEntry();
    out.close();
    DefaultMessage message = new DefaultMessage();
    Map<String, Object> newHeaders = new CaseInsensitiveMap(headers);
    newHeaders.put("entryName", entryName);
    newHeaders.put("unzippedFiles", unzippedFiles);
    message.setHeaders(newHeaders);
    message.setBody(out.toByteArray());
    answer.add(message);
   }
   zis.close();
  } catch (Throwable e) {
   e.printStackTrace();
  }
  return answer;
 }
}

Enfin, le code utilisé pour créer des archives zip utilisables (via CustomizedZipDataFormat) provient de cette page.

J'espère que tout ce code ne paraît pas trop indigeste, mais en y regardant de plus près, on s'aperçoit que Camel permet assez facilement de mettre en place ce genre de flux de données, en peu de lignes de code et surtout de manière plutôt lisible. La documentation est d'ailleurs incroyablement bien faite pour ce qui concerne la description des patrons et des composants natifs.

Le seul bémol que j'ajouterai est la gestion des formats ZIP. En effet, par défaut, Camel crée des DeflaterOutputStream : je ne sais pas trop d'où provient ce format, mais en tout cas il ne permet pas directement de créer des archives lisibles. Il faut donc explicitement les convertir en ZipOutputStream classique.

N'hésitez pas à m'indiquer si vous avez déjà utilisé cet outil et surtout si vous voyez des façons d'améliorer ce que j'ai présenté!

Sources :


Fichier(s) joint(s) :